第一章.WaterMark
1.WaterMark API
一.升序API
为什么数据没有乱序的情况下还要使用WaterMark
因为只要程序的时间语义是事件时间,就必须使用WaterMark
升序条件下,WaterMark = EventTime(截止目前最大事件时间) - 1(毫秒)
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.util.Arrays;public class $01_WaterMarkMoNo {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {@Overridepublic long extractTimestamp(WaterSensor element, long recordTimestamp) {return element.getTs() * 1000L;}}));WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10)));SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();env.execute();}}
二.乱序API
窗口触发条件:
WaterMark = EventTime(截止目前最大事件时间) - maxOutOfOrderness(等待时间) - 1(毫秒)
WaterMark >= maxTimestamp(窗口最大时间戳)
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.time.Duration;import java.util.Arrays;public class $02_WaterMarkOutOfOrderness {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy/*** 乱序程度: 1.靠经验值: 2.靠抽样估算* 生产环境: 秒级或者分钟级*/.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4)).withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {@Overridepublic long extractTimestamp(WaterSensor element, long recordTimestamp) {return element.getTs() * 1000L;}}));WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10)));SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();env.execute();}}
2.自定义WaterMark
有两种风格的WaterMark生产方式:periodic(周期性)和punctuated(间歇性),都需要继承接口:WatermarkGenerztor
一.周期性
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.*;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.util.Arrays;/*** WaterMark自定义------>周期性生成*/public class $03_Mygenerator {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//TODO 指定watermark生成周期,默认是200msenv.getConfig().setAutoWatermarkInterval(2000L);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy/*** 乱序程度: 1.靠经验值: 2.靠抽样估算* 生产环境: 秒级或者分钟级*/.forGenerator(new WatermarkGeneratorSupplier<WaterSensor>() {@Overridepublic WatermarkGenerator<WaterSensor> createWatermarkGenerator(Context context) {return new MyPeriodicGenerator<WaterSensor>(3000L);}}).withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {@Overridepublic long extractTimestamp(WaterSensor element, long recordTimestamp) {return element.getTs() * 1000L;}}));WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10)));SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();env.execute();}public static class MyPeriodicGenerator<T> implements WatermarkGenerator<T>{long maxOutOfOrderness;long maxTs;public MyPeriodicGenerator(long maxOutOfOrderness) {this.maxOutOfOrderness = maxOutOfOrderness;this.maxTs = Long.MIN_VALUE + maxOutOfOrderness;}/*** 每来一次数据,调用一次该方法* @param event* @param eventTimestamp* @param output*/@Overridepublic void onEvent(T event, long eventTimestamp, WatermarkOutput output) {System.out.println("onEvent....");maxTs = Math.max(maxTs,eventTimestamp);}/*** 每来一个固定周期,调用一次该方法* @param output*/@Overridepublic void onPeriodicEmit(WatermarkOutput output) {System.out.println("onPeriodicEmit");output.emitWatermark(new Watermark(maxTs - maxOutOfOrderness));}}}
二.间歇性
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.*;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.util.Arrays;/*** WaterMark自定义------>间歇性生成*/public class $04_MygeneratorByPuntucate {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//TODO 指定watermark生成周期,默认是200msenv.getConfig().setAutoWatermarkInterval(2000L);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy/*** 乱序程度: 1.靠经验值: 2.靠抽样估算* 生产环境: 秒级或者分钟级*/.forGenerator(new WatermarkGeneratorSupplier<WaterSensor>() {@Overridepublic WatermarkGenerator<WaterSensor> createWatermarkGenerator(Context context) {return new MyPunctuateGenerator<WaterSensor>(3000L);}}).withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {@Overridepublic long extractTimestamp(WaterSensor element, long recordTimestamp) {return element.getTs() * 1000L;}}));WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10)));SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();env.execute();}public static class MyPunctuateGenerator<T> implements WatermarkGenerator<T>{long maxOutOfOrderness;long maxTs;public MyPunctuateGenerator(long maxOutOfOrderness) {this.maxOutOfOrderness = maxOutOfOrderness;this.maxTs = Long.MIN_VALUE + maxOutOfOrderness;}/*** 每来一次数据,调用一次该方法* @param event* @param eventTimestamp* @param output*/@Overridepublic void onEvent(T event, long eventTimestamp, WatermarkOutput output) {System.out.println("onEvent....");maxTs = Math.max(maxTs,eventTimestamp);Watermark watermark = new Watermark(maxTs - maxOutOfOrderness);System.out.println("onEvent>>>watermark=" + watermark.getTimestamp());output.emitWatermark(watermark);}/*** 每来一个固定周期,调用一次该方法* @param output*/@Overridepublic void onPeriodicEmit(WatermarkOutput output) {System.out.println("onPeriodicEmit");}}}
3.WaterMark老版本写法
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.*;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.timestamps.AscendingTimestampExtractor;import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.util.Arrays;/*** WaterMark老版本写法*/public class $05_WaterMarkOldVersion {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}})/*** 1.10及以前版本的写法**/.assignTimestampsAndWatermarks(/*** TODO 升序写法*//*new AscendingTimestampExtractor<WaterSensor>() {@Overridepublic long extractAscendingTimestamp(WaterSensor element) {return element.getTs() * 1000L;}}*//*** TODO 降序写法*/new BoundedOutOfOrdernessTimestampExtractor<WaterSensor>(Time.seconds(3)) {@Overridepublic long extractTimestamp(WaterSensor element) {return element.getTs() * 1000L;}});WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10)));SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();env.execute();}}
4.多并行度下WaterMark的传递
多并行度的条件下,向下游传递WaterMark的时候,总是以最小的那个WaterMark为准,木桶原理
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.time.Duration;import java.util.Arrays;/*** 多并行度下WaterMark的分析*/public class $06_WaterMarkMul {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(2);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4)).withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {@Overridepublic long extractTimestamp(WaterSensor element, long recordTimestamp) {return element.getTs() * 1000L;}}));WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10)));SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();env.execute();}}
![day05[Flink流处理高阶编程(中)] - 图1](/uploads/projects/liuye-6lcqc@ddtw8t/3f47a5f0f563a36dbe64d12b90740f61.png)
![day05[Flink流处理高阶编程(中)] - 图2](/uploads/projects/liuye-6lcqc@ddtw8t/77c7efae585a9fdc61b949cf56d4dfda.png)
5.迟到数据
一.窗口允许迟到
- 当时间进展 >= maxTs,正常触发输出
- 当maxTs < 时间进展 < maxTs + 允许迟到时间,每来一条数据,都会触发一次计算输出
- 当时间进展 >= maxTs + 允许迟到时间,窗口关闭,再有迟到的数据来,不会处理
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import java.time.Duration;import java.util.Arrays;public class $07_WaterMarkAllowedLateness {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4)).withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {@Overridepublic long extractTimestamp(WaterSensor element, long recordTimestamp) {return element.getTs() * 1000L;}}));WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10))).allowedLateness(Time.seconds(3));SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();env.execute();}}
二.侧输出流
接收窗口关闭之后的迟到数据
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.datastream.WindowedStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.util.Collector;import org.apache.flink.util.OutputTag;import java.time.Duration;import java.util.Arrays;public class $08_WaterMarkSideOutput {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4)).withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {@Overridepublic long extractTimestamp(WaterSensor element, long recordTimestamp) {return element.getTs() * 1000L;}}));/*** Tag要用匿名内部类写法* 1.指定泛型* 2.加大括号*/OutputTag<WaterSensor> lateTag = new OutputTag<WaterSensor>("lateTag") {};WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.keyBy(sensor -> sensor.getId()).window(TumblingEventTimeWindows.of(Time.seconds(10))).sideOutputLateData(lateTag);SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {out.collect("key" + s + "\n"+ "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"+ "数据为:" + Arrays.asList(elements).toString() + "\n"+ "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");}});resultDS.print();//获取侧输出流的数据resultDS.getSideOutput(lateTag).print("late");env.execute();}}
6.WaterMark总结
- WaterMark理解
- 用来衡量事件时间的进展
- 用来出来乱序问题
- 是一个特殊的时间戳,从生成位置插入到流里面
- 单调不减的
- 用来触发窗口等
- Flink认为,事件时间小于WaterMark的数据应该都已经处理过了,如果后续还有事件时间小于WaterMark的数据来,称为迟到数据
- WaterMark的生成方式
- 周期性(onPeriodicEmit里发射WaterMark):默认200ms
- 间歇性(onEvent里发射WaterMark)
- WaterMark的写法,生成逻辑
- 升序写法: forMonotonousTimestamps (WaterMark = maxTs - 1ms)
- 乱序写法: forBoundedOutOfOrderness(Duration) (WaterMark = maxTs - Duration - 1ms)
- WaterMark的传递
- 从指定位置插入的,向下游传递
- 一对多(发送者):广播
- 多对一(接收者):取最小(木桶原理)
- 多对多:上面两者的结合,拆开来看
- Flink对乱序和迟到数据的处理方式
- WaterMark处理乱序
- 窗口允许迟到
- 侧输出流
第二章.ProcessFunction API
之前的转换算子是无法访问事件的时间戳信息,和水位线信息的,ProcessFunction则可以用来构建事件驱动的应用以及实现自定义的业务逻辑
1.Process(时间戳)
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.ProcessFunction;import org.apache.flink.util.Collector;/*** 上下文.timestamp取的是事件时间,要求我们要指定时间的提取,否则为null*/public class $09_ProcessDemo {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((data,ts)-> data.getTs() * 1000L));sensorDS.process(new ProcessFunction<WaterSensor, String>() {@Overridepublic void processElement(WaterSensor value, Context context, Collector<String> collector) throws Exception {collector.collect("数据为:" + value.toString() + "\n" +"timestamp=" + context.timestamp() + "\n\n");}}).print();env.execute();}}
2.Process(侧输出流)
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.datastream.KeyedStream;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.KeyedProcessFunction;import org.apache.flink.streaming.api.functions.ProcessFunction;import org.apache.flink.util.Collector;import org.apache.flink.util.OutputTag;/*** 侧输出流可以报警,实现分流*/public class $09_ProcessSideOutput {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((data,ts)-> data.getTs() * 1000L));KeyedStream<WaterSensor, String> sensorKS = sensorDS.keyBy(sensor -> sensor.getId());SingleOutputStreamOperator<String> resultDS = sensorKS.process(new KeyedProcessFunction<String, WaterSensor, String>() {@Overridepublic void processElement(WaterSensor waterSensor, Context context, Collector<String> collector) throws Exception {OutputTag<Integer> outputTag = new OutputTag<Integer>("baojing") {};if (waterSensor.getVc() > 10) {context.output(outputTag, waterSensor.getVc());} else {collector.collect(waterSensor.toString());}}});resultDS.print();DataStream<Integer> sideOutput = resultDS.getSideOutput(new OutputTag<Integer>("baojing") {});sideOutput.print("side");env.execute();}}
第三章.定时器
基于处理时间或事件时间处理过一个元素之后,注册一个定时器,然后在指定的时间执行
1.基于处理时间的定时器
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.TimerService;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.datastream.KeyedStream;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.KeyedProcessFunction;import org.apache.flink.util.Collector;import org.apache.flink.util.OutputTag;/*** 基于处理时间的定时器*/public class $11_TimeServiceDemo1 {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((data,ts)-> data.getTs() * 1000L));KeyedStream<WaterSensor, String> sensorKS = sensorDS.keyBy(sensor -> sensor.getId());SingleOutputStreamOperator<String> resultDS = sensorKS.process(new KeyedProcessFunction<String, WaterSensor, String>() {@Overridepublic void processElement(WaterSensor waterSensor, Context context, Collector<String> collector) throws Exception {TimerService timerService = context.timerService();System.out.println("currentProcessingTime=" + timerService.currentProcessingTime() + "\n"+ "currentWatermark=" + timerService.currentWatermark());timerService.registerProcessingTimeTimer(System.currentTimeMillis() + 1000L);}/*** 定时器触发,当时间进展到了注册的时间,调用该方法* @param timestamp* @param ctx* @param out* @throws Exception*/@Overridepublic void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {System.out.println("onTimer的timestamp=" + timestamp);}});resultDS.print();env.execute();}}
2.基于事件时间的定时器
package com.atguigu.flink.day05;import com.atguigu.flink.day02.pojo.WaterSensor;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.streaming.api.TimerService;import org.apache.flink.streaming.api.datastream.KeyedStream;import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.KeyedProcessFunction;import org.apache.flink.util.Collector;/*** 基于处理时间的定时器*/public class $12_TimeServiceDemo2 {public static void main(String[] args) throws Exception {StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999).map(new MapFunction<String, WaterSensor>() {@Overridepublic WaterSensor map(String value) throws Exception {String[] line = value.split(",");return new WaterSensor(line[0],Long.parseLong(line[1]),Integer.parseInt(line[2]));}}).assignTimestampsAndWatermarks(WatermarkStrategy.<WaterSensor>forMonotonousTimestamps().withTimestampAssigner((data,ts)-> data.getTs() * 1000L));KeyedStream<WaterSensor, String> sensorKS = sensorDS.keyBy(sensor -> sensor.getId());SingleOutputStreamOperator<String> resultDS = sensorKS.process(new KeyedProcessFunction<String, WaterSensor, String>() {@Overridepublic void processElement(WaterSensor waterSensor, Context context, Collector<String> collector) throws Exception {TimerService timerService = context.timerService();System.out.println("currentProcessingTime=" + timerService.currentProcessingTime() + "\n"+ "currentWatermark=" + timerService.currentWatermark());timerService.registerEventTimeTimer(context.timestamp() + 5000L);}/*** 定时器触发,当时间进展到了注册的时间,调用该方法* @param timestamp* @param ctx* @param out* @throws Exception*/@Overridepublic void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {System.out.println("onTimer的timestamp=" + timestamp +"\n"+ "watermark=" + ctx.timerService().currentWatermark());}});resultDS.print();env.execute();}}
启动socket测试
[atguigu@hadoop102 ~]$ nc -lk 9999s_1,1,1s_1,6,1s_1,7,1
![day05[Flink流处理高阶编程(中)] - 图3](/uploads/projects/liuye-6lcqc@ddtw8t/8913ea1549b44581207aa328d3a7d1c8.png)
