前言
【Flink 监控状态案例之温度跳变发警告】在项目中,经常遇到这样的场景,对于一批源源不断进入flink的数据源,需要检测某种类型的数据连续两次之间的数值变化范围,如果这个变化的值大于或者小于一定的标准值,将给出相应的告警;
在上一篇关于flink的常用状态管理的总结文章中,我们了解到了flink的常用的几种状态,如果应对这个场景,该使用哪种状态管理比较好呢?很明显是键控状态了;
通常来说,在实际的业务数据流中,都会有一些唯一标识数据的字段,那么通过这个字段做keyby的操作,接下来就可以使用键控状态做处理了;
下面看具体的代码实现:
import com.congge.source.SensorReading;import org.apache.flink.api.common.functions.RichFlatMapFunction;import org.apache.flink.api.common.state.ValueState;import org.apache.flink.api.common.state.ValueStateDescriptor;import org.apache.flink.api.java.tuple.Tuple3;import org.apache.flink.configuration.Configuration;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.util.Collector;public class KeyedStateApplicationCase1 {public static void main(String[] args) throws Exception{StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);// socket文本流DataStream inputStream = env.socketTextStream("localhost", 7777);// 转换成SensorReading类型DataStream dataStream = inputStream.map(line -> {String[] fields = line.split(",");return new SensorReading(fields[0], new Long(fields[1]), new Double(fields[2]));});// 定义一个flatmap操作,检测温度跳变,输出报警SingleOutputStreamOperator resultStream = dataStream.keyBy("id").flatMap(new TempChangeWarning(10.0));resultStream.print();env.execute();}public static class TempChangeWarning extends RichFlatMapFunction>{// 私有属性,温度跳变阈值private Double threshold;public TempChangeWarning(Double threshold) {this.threshold = threshold;}// 定义状态,保存上一次的温度值private ValueState
需要重点说明一点的是,这里keyby之后使用的是flatmap,因为flatmap会将过来的数据进行扁平化处理,同时由于需要记录数据的上下文状态,使用了RichFlatMapFunction;
- 春季老年人吃什么养肝?土豆、米饭换着吃
- 三八妇女节节日祝福分享 三八妇女节节日语录
- 老人谨慎!选好你的“第三只脚”
- 校方进行了深刻的反思 青岛一大学生坠亡校方整改校规
- 脸皮厚的人长寿!有这特征的老人最长寿
- 长寿秘诀:记住这10大妙招 100%增寿
- 春季老年人心血管病高发 3条保命要诀
- 眼睛花不花要看四十八 老年人怎样延缓老花眼
- 香槟然能防治老年痴呆症? 一天三杯它人到90不痴呆
- 老人手抖的原因 为什么老人手会抖
