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)
还没有任何评论哟~
