Advertisement

Flink 检测温度持续上升趋势并触发警报输出

阅读量:

需求:当Flink检测到温度在10秒内持续上升时,系统应触发报警机制

方案:为实现该功能,我们采用了keyBy函数,原因在于唯有KeyedProcessFunction能够处理KeyedStream类型的数据流

以下将简要说明KeyedStream的相关概念:

KeyedStream的父类为RichFunction,其在数据分流后会对每个元素执行一次KeyedProcessFunction中的elementProcess方法。通过Context对象,可以调用时间服务、注册定时器,并获取当前水位线、处理时间等关键信息。

思路:

复制代码

2、在processElement方法中,每次对数据进行处理时,会获取前一次记录的温度数值及对应状态。当检测到温度呈现上升趋势且尚未设置定时器时,将启动一个延迟10秒的定时机制,以等待后续onTimer事件的触发。

3、一旦发现温度出现下降情况,则立即取消已设置的定时器。

复制代码
  
    
 import com.atguigu.apitest.beans.SensorReading;
    
 import org.apache.flink.api.common.state.ValueState;
    
 import org.apache.flink.api.common.state.Valu

全部评论 (0)

还没有任何评论哟~