第一章.WaterMark

1.WaterMark API

一.升序API

为什么数据没有乱序的情况下还要使用WaterMark

因为只要程序的时间语义是事件时间,就必须使用WaterMark

升序条件下,WaterMark = EventTime(截止目前最大事件时间) - 1(毫秒)

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
  4. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  5. import org.apache.flink.api.common.functions.MapFunction;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.datastream.WindowedStream;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  10. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  11. import org.apache.flink.streaming.api.windowing.time.Time;
  12. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  13. import org.apache.flink.util.Collector;
  14. import java.util.Arrays;
  15. public class $01_WaterMarkMoNo {
  16. public static void main(String[] args) throws Exception {
  17. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.setParallelism(1);
  19. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  20. .map(new MapFunction<String, WaterSensor>() {
  21. @Override
  22. public WaterSensor map(String value) throws Exception {
  23. String[] line = value.split(",");
  24. return new WaterSensor(
  25. line[0],
  26. Long.parseLong(line[1]),
  27. Integer.parseInt(line[2])
  28. );
  29. }
  30. })
  31. .assignTimestampsAndWatermarks(
  32. WatermarkStrategy
  33. .<WaterSensor>forMonotonousTimestamps()
  34. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  35. @Override
  36. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  37. return element.getTs() * 1000L;
  38. }
  39. })
  40. );
  41. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  42. keyBy(sensor -> sensor.getId())
  43. .window(TumblingEventTimeWindows.of(Time.seconds(10)));
  44. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  45. @Override
  46. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  47. out.collect("key" + s + "\n"
  48. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  49. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  50. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  51. }
  52. });
  53. resultDS.print();
  54. env.execute();
  55. }
  56. }

二.乱序API

窗口触发条件:

WaterMark = EventTime(截止目前最大事件时间) - maxOutOfOrderness(等待时间) - 1(毫秒)

WaterMark >= maxTimestamp(窗口最大时间戳)

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
  4. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  5. import org.apache.flink.api.common.functions.MapFunction;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.datastream.WindowedStream;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  10. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  11. import org.apache.flink.streaming.api.windowing.time.Time;
  12. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  13. import org.apache.flink.util.Collector;
  14. import java.time.Duration;
  15. import java.util.Arrays;
  16. public class $02_WaterMarkOutOfOrderness {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  21. .map(new MapFunction<String, WaterSensor>() {
  22. @Override
  23. public WaterSensor map(String value) throws Exception {
  24. String[] line = value.split(",");
  25. return new WaterSensor(
  26. line[0],
  27. Long.parseLong(line[1]),
  28. Integer.parseInt(line[2])
  29. );
  30. }
  31. })
  32. .assignTimestampsAndWatermarks(
  33. WatermarkStrategy
  34. /**
  35. * 乱序程度: 1.靠经验值: 2.靠抽样估算
  36. * 生产环境: 秒级或者分钟级
  37. */
  38. .<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4))
  39. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  40. @Override
  41. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  42. return element.getTs() * 1000L;
  43. }
  44. })
  45. );
  46. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  47. keyBy(sensor -> sensor.getId())
  48. .window(TumblingEventTimeWindows.of(Time.seconds(10)));
  49. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  50. @Override
  51. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  52. out.collect("key" + s + "\n"
  53. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  54. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  55. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  56. }
  57. });
  58. resultDS.print();
  59. env.execute();
  60. }
  61. }

2.自定义WaterMark

有两种风格的WaterMark生产方式:periodic(周期性)和punctuated(间歇性),都需要继承接口:WatermarkGenerztor

一.周期性

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.*;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  6. import org.apache.flink.streaming.api.datastream.WindowedStream;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  9. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  10. import org.apache.flink.streaming.api.windowing.time.Time;
  11. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  12. import org.apache.flink.util.Collector;
  13. import java.util.Arrays;
  14. /**
  15. * WaterMark自定义------>周期性生成
  16. */
  17. public class $03_Mygenerator {
  18. public static void main(String[] args) throws Exception {
  19. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  20. env.setParallelism(1);
  21. //TODO 指定watermark生成周期,默认是200ms
  22. env.getConfig().setAutoWatermarkInterval(2000L);
  23. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  24. .map(new MapFunction<String, WaterSensor>() {
  25. @Override
  26. public WaterSensor map(String value) throws Exception {
  27. String[] line = value.split(",");
  28. return new WaterSensor(
  29. line[0],
  30. Long.parseLong(line[1]),
  31. Integer.parseInt(line[2])
  32. );
  33. }
  34. })
  35. .assignTimestampsAndWatermarks(
  36. WatermarkStrategy
  37. /**
  38. * 乱序程度: 1.靠经验值: 2.靠抽样估算
  39. * 生产环境: 秒级或者分钟级
  40. */
  41. .forGenerator(new WatermarkGeneratorSupplier<WaterSensor>() {
  42. @Override
  43. public WatermarkGenerator<WaterSensor> createWatermarkGenerator(Context context) {
  44. return new MyPeriodicGenerator<WaterSensor>(3000L);
  45. }
  46. })
  47. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  48. @Override
  49. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  50. return element.getTs() * 1000L;
  51. }
  52. })
  53. );
  54. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  55. keyBy(sensor -> sensor.getId())
  56. .window(TumblingEventTimeWindows.of(Time.seconds(10)));
  57. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  58. @Override
  59. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  60. out.collect("key" + s + "\n"
  61. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  62. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  63. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  64. }
  65. });
  66. resultDS.print();
  67. env.execute();
  68. }
  69. public static class MyPeriodicGenerator<T> implements WatermarkGenerator<T>{
  70. long maxOutOfOrderness;
  71. long maxTs;
  72. public MyPeriodicGenerator(long maxOutOfOrderness) {
  73. this.maxOutOfOrderness = maxOutOfOrderness;
  74. this.maxTs = Long.MIN_VALUE + maxOutOfOrderness;
  75. }
  76. /**
  77. * 每来一次数据,调用一次该方法
  78. * @param event
  79. * @param eventTimestamp
  80. * @param output
  81. */
  82. @Override
  83. public void onEvent(T event, long eventTimestamp, WatermarkOutput output) {
  84. System.out.println("onEvent....");
  85. maxTs = Math.max(maxTs,eventTimestamp);
  86. }
  87. /**
  88. * 每来一个固定周期,调用一次该方法
  89. * @param output
  90. */
  91. @Override
  92. public void onPeriodicEmit(WatermarkOutput output) {
  93. System.out.println("onPeriodicEmit");
  94. output.emitWatermark(new Watermark(maxTs - maxOutOfOrderness));
  95. }
  96. }
  97. }

二.间歇性

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.*;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  6. import org.apache.flink.streaming.api.datastream.WindowedStream;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  9. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  10. import org.apache.flink.streaming.api.windowing.time.Time;
  11. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  12. import org.apache.flink.util.Collector;
  13. import java.util.Arrays;
  14. /**
  15. * WaterMark自定义------>间歇性生成
  16. */
  17. public class $04_MygeneratorByPuntucate {
  18. public static void main(String[] args) throws Exception {
  19. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  20. env.setParallelism(1);
  21. //TODO 指定watermark生成周期,默认是200ms
  22. env.getConfig().setAutoWatermarkInterval(2000L);
  23. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  24. .map(new MapFunction<String, WaterSensor>() {
  25. @Override
  26. public WaterSensor map(String value) throws Exception {
  27. String[] line = value.split(",");
  28. return new WaterSensor(
  29. line[0],
  30. Long.parseLong(line[1]),
  31. Integer.parseInt(line[2])
  32. );
  33. }
  34. })
  35. .assignTimestampsAndWatermarks(
  36. WatermarkStrategy
  37. /**
  38. * 乱序程度: 1.靠经验值: 2.靠抽样估算
  39. * 生产环境: 秒级或者分钟级
  40. */
  41. .forGenerator(new WatermarkGeneratorSupplier<WaterSensor>() {
  42. @Override
  43. public WatermarkGenerator<WaterSensor> createWatermarkGenerator(Context context) {
  44. return new MyPunctuateGenerator<WaterSensor>(3000L);
  45. }
  46. })
  47. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  48. @Override
  49. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  50. return element.getTs() * 1000L;
  51. }
  52. })
  53. );
  54. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  55. keyBy(sensor -> sensor.getId())
  56. .window(TumblingEventTimeWindows.of(Time.seconds(10)));
  57. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  58. @Override
  59. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  60. out.collect("key" + s + "\n"
  61. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  62. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  63. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  64. }
  65. });
  66. resultDS.print();
  67. env.execute();
  68. }
  69. public static class MyPunctuateGenerator<T> implements WatermarkGenerator<T>{
  70. long maxOutOfOrderness;
  71. long maxTs;
  72. public MyPunctuateGenerator(long maxOutOfOrderness) {
  73. this.maxOutOfOrderness = maxOutOfOrderness;
  74. this.maxTs = Long.MIN_VALUE + maxOutOfOrderness;
  75. }
  76. /**
  77. * 每来一次数据,调用一次该方法
  78. * @param event
  79. * @param eventTimestamp
  80. * @param output
  81. */
  82. @Override
  83. public void onEvent(T event, long eventTimestamp, WatermarkOutput output) {
  84. System.out.println("onEvent....");
  85. maxTs = Math.max(maxTs,eventTimestamp);
  86. Watermark watermark = new Watermark(maxTs - maxOutOfOrderness);
  87. System.out.println("onEvent>>>watermark=" + watermark.getTimestamp());
  88. output.emitWatermark(watermark);
  89. }
  90. /**
  91. * 每来一个固定周期,调用一次该方法
  92. * @param output
  93. */
  94. @Override
  95. public void onPeriodicEmit(WatermarkOutput output) {
  96. System.out.println("onPeriodicEmit");
  97. }
  98. }
  99. }

3.WaterMark老版本写法

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.*;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  6. import org.apache.flink.streaming.api.datastream.WindowedStream;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.timestamps.AscendingTimestampExtractor;
  9. import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor;
  10. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  11. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  12. import org.apache.flink.streaming.api.windowing.time.Time;
  13. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  14. import org.apache.flink.util.Collector;
  15. import java.util.Arrays;
  16. /**
  17. * WaterMark老版本写法
  18. */
  19. public class $05_WaterMarkOldVersion {
  20. public static void main(String[] args) throws Exception {
  21. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  22. env.setParallelism(1);
  23. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  24. .map(new MapFunction<String, WaterSensor>() {
  25. @Override
  26. public WaterSensor map(String value) throws Exception {
  27. String[] line = value.split(",");
  28. return new WaterSensor(
  29. line[0],
  30. Long.parseLong(line[1]),
  31. Integer.parseInt(line[2])
  32. );
  33. }
  34. })
  35. /**
  36. * 1.10及以前版本的写法
  37. *
  38. */
  39. .assignTimestampsAndWatermarks(
  40. /**
  41. * TODO 升序写法
  42. */
  43. /*new AscendingTimestampExtractor<WaterSensor>() {
  44. @Override
  45. public long extractAscendingTimestamp(WaterSensor element) {
  46. return element.getTs() * 1000L;
  47. }
  48. }*/
  49. /**
  50. * TODO 降序写法
  51. */
  52. new BoundedOutOfOrdernessTimestampExtractor<WaterSensor>(Time.seconds(3)) {
  53. @Override
  54. public long extractTimestamp(WaterSensor element) {
  55. return element.getTs() * 1000L;
  56. }
  57. }
  58. );
  59. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  60. keyBy(sensor -> sensor.getId())
  61. .window(TumblingEventTimeWindows.of(Time.seconds(10)));
  62. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  63. @Override
  64. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  65. out.collect("key" + s + "\n"
  66. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  67. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  68. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  69. }
  70. });
  71. resultDS.print();
  72. env.execute();
  73. }
  74. }

4.多并行度下WaterMark的传递

多并行度的条件下,向下游传递WaterMark的时候,总是以最小的那个WaterMark为准,木桶原理

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
  4. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  5. import org.apache.flink.api.common.functions.MapFunction;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.datastream.WindowedStream;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  10. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  11. import org.apache.flink.streaming.api.windowing.time.Time;
  12. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  13. import org.apache.flink.util.Collector;
  14. import java.time.Duration;
  15. import java.util.Arrays;
  16. /**
  17. * 多并行度下WaterMark的分析
  18. */
  19. public class $06_WaterMarkMul {
  20. public static void main(String[] args) throws Exception {
  21. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  22. env.setParallelism(2);
  23. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  24. .map(new MapFunction<String, WaterSensor>() {
  25. @Override
  26. public WaterSensor map(String value) throws Exception {
  27. String[] line = value.split(",");
  28. return new WaterSensor(
  29. line[0],
  30. Long.parseLong(line[1]),
  31. Integer.parseInt(line[2])
  32. );
  33. }
  34. })
  35. .assignTimestampsAndWatermarks(
  36. WatermarkStrategy
  37. .<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4))
  38. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  39. @Override
  40. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  41. return element.getTs() * 1000L;
  42. }
  43. })
  44. );
  45. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  46. keyBy(sensor -> sensor.getId())
  47. .window(TumblingEventTimeWindows.of(Time.seconds(10)));
  48. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  49. @Override
  50. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  51. out.collect("key" + s + "\n"
  52. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  53. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  54. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  55. }
  56. });
  57. resultDS.print();
  58. env.execute();
  59. }
  60. }

day05[Flink流处理高阶编程(中)] - 图1

day05[Flink流处理高阶编程(中)] - 图2

5.迟到数据

一.窗口允许迟到

  1. 当时间进展 >= maxTs,正常触发输出
  2. 当maxTs < 时间进展 < maxTs + 允许迟到时间,每来一条数据,都会触发一次计算输出
  3. 当时间进展 >= maxTs + 允许迟到时间,窗口关闭,再有迟到的数据来,不会处理
  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
  4. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  5. import org.apache.flink.api.common.functions.MapFunction;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.datastream.WindowedStream;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  10. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  11. import org.apache.flink.streaming.api.windowing.time.Time;
  12. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  13. import org.apache.flink.util.Collector;
  14. import java.time.Duration;
  15. import java.util.Arrays;
  16. public class $07_WaterMarkAllowedLateness {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  21. .map(new MapFunction<String, WaterSensor>() {
  22. @Override
  23. public WaterSensor map(String value) throws Exception {
  24. String[] line = value.split(",");
  25. return new WaterSensor(
  26. line[0],
  27. Long.parseLong(line[1]),
  28. Integer.parseInt(line[2])
  29. );
  30. }
  31. })
  32. .assignTimestampsAndWatermarks(
  33. WatermarkStrategy
  34. .<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4))
  35. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  36. @Override
  37. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  38. return element.getTs() * 1000L;
  39. }
  40. })
  41. );
  42. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  43. keyBy(sensor -> sensor.getId())
  44. .window(TumblingEventTimeWindows.of(Time.seconds(10)))
  45. .allowedLateness(Time.seconds(3));
  46. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  47. @Override
  48. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  49. out.collect("key" + s + "\n"
  50. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  51. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  52. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  53. }
  54. });
  55. resultDS.print();
  56. env.execute();
  57. }
  58. }

二.侧输出流

接收窗口关闭之后的迟到数据

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
  4. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  5. import org.apache.flink.api.common.functions.MapFunction;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.datastream.WindowedStream;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  10. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  11. import org.apache.flink.streaming.api.windowing.time.Time;
  12. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  13. import org.apache.flink.util.Collector;
  14. import org.apache.flink.util.OutputTag;
  15. import java.time.Duration;
  16. import java.util.Arrays;
  17. public class $08_WaterMarkSideOutput {
  18. public static void main(String[] args) throws Exception {
  19. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  20. env.setParallelism(1);
  21. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  22. .map(new MapFunction<String, WaterSensor>() {
  23. @Override
  24. public WaterSensor map(String value) throws Exception {
  25. String[] line = value.split(",");
  26. return new WaterSensor(
  27. line[0],
  28. Long.parseLong(line[1]),
  29. Integer.parseInt(line[2])
  30. );
  31. }
  32. })
  33. .assignTimestampsAndWatermarks(
  34. WatermarkStrategy
  35. .<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(4))
  36. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  37. @Override
  38. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  39. return element.getTs() * 1000L;
  40. }
  41. })
  42. );
  43. /**
  44. * Tag要用匿名内部类写法
  45. * 1.指定泛型
  46. * 2.加大括号
  47. */
  48. OutputTag<WaterSensor> lateTag = new OutputTag<WaterSensor>("lateTag") {};
  49. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  50. keyBy(sensor -> sensor.getId())
  51. .window(TumblingEventTimeWindows.of(Time.seconds(10)))
  52. .sideOutputLateData(lateTag);
  53. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  54. @Override
  55. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  56. out.collect("key" + s + "\n"
  57. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  58. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  59. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  60. }
  61. });
  62. resultDS.print();
  63. //获取侧输出流的数据
  64. resultDS.getSideOutput(lateTag).print("late");
  65. env.execute();
  66. }
  67. }

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(时间戳)

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.functions.ProcessFunction;
  8. import org.apache.flink.util.Collector;
  9. /**
  10. * 上下文.timestamp取的是事件时间,要求我们要指定时间的提取,否则为null
  11. */
  12. public class $09_ProcessDemo {
  13. public static void main(String[] args) throws Exception {
  14. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  15. env.setParallelism(1);
  16. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  17. .map(new MapFunction<String, WaterSensor>() {
  18. @Override
  19. public WaterSensor map(String value) throws Exception {
  20. String[] line = value.split(",");
  21. return new WaterSensor(
  22. line[0],
  23. Long.parseLong(line[1]),
  24. Integer.parseInt(line[2])
  25. );
  26. }
  27. })
  28. .assignTimestampsAndWatermarks(
  29. WatermarkStrategy
  30. .<WaterSensor>forMonotonousTimestamps()
  31. .withTimestampAssigner((data,ts)-> data.getTs() * 1000L)
  32. );
  33. sensorDS
  34. .process(
  35. new ProcessFunction<WaterSensor, String>() {
  36. @Override
  37. public void processElement(WaterSensor value, Context context, Collector<String> collector) throws Exception {
  38. collector.collect("数据为:" + value.toString() + "\n" +
  39. "timestamp=" + context.timestamp() + "\n\n");
  40. }
  41. }
  42. )
  43. .print();
  44. env.execute();
  45. }
  46. }

2.Process(侧输出流)

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.datastream.DataStream;
  6. import org.apache.flink.streaming.api.datastream.KeyedStream;
  7. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
  10. import org.apache.flink.streaming.api.functions.ProcessFunction;
  11. import org.apache.flink.util.Collector;
  12. import org.apache.flink.util.OutputTag;
  13. /**
  14. * 侧输出流可以报警,实现分流
  15. */
  16. public class $09_ProcessSideOutput {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  21. .map(new MapFunction<String, WaterSensor>() {
  22. @Override
  23. public WaterSensor map(String value) throws Exception {
  24. String[] line = value.split(",");
  25. return new WaterSensor(
  26. line[0],
  27. Long.parseLong(line[1]),
  28. Integer.parseInt(line[2])
  29. );
  30. }
  31. })
  32. .assignTimestampsAndWatermarks(
  33. WatermarkStrategy
  34. .<WaterSensor>forMonotonousTimestamps()
  35. .withTimestampAssigner((data,ts)-> data.getTs() * 1000L)
  36. );
  37. KeyedStream<WaterSensor, String> sensorKS = sensorDS.keyBy(sensor -> sensor.getId());
  38. SingleOutputStreamOperator<String> resultDS = sensorKS.process(new KeyedProcessFunction<String, WaterSensor, String>() {
  39. @Override
  40. public void processElement(WaterSensor waterSensor, Context context, Collector<String> collector) throws Exception {
  41. OutputTag<Integer> outputTag = new OutputTag<Integer>("baojing") {
  42. };
  43. if (waterSensor.getVc() > 10) {
  44. context.output(outputTag, waterSensor.getVc());
  45. } else {
  46. collector.collect(waterSensor.toString());
  47. }
  48. }
  49. });
  50. resultDS.print();
  51. DataStream<Integer> sideOutput = resultDS.getSideOutput(new OutputTag<Integer>("baojing") {
  52. });
  53. sideOutput.print("side");
  54. env.execute();
  55. }
  56. }

第三章.定时器

基于处理时间或事件时间处理过一个元素之后,注册一个定时器,然后在指定的时间执行

1.基于处理时间的定时器

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.TimerService;
  6. import org.apache.flink.streaming.api.datastream.DataStream;
  7. import org.apache.flink.streaming.api.datastream.KeyedStream;
  8. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  9. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  10. import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
  11. import org.apache.flink.util.Collector;
  12. import org.apache.flink.util.OutputTag;
  13. /**
  14. * 基于处理时间的定时器
  15. */
  16. public class $11_TimeServiceDemo1 {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  21. .map(new MapFunction<String, WaterSensor>() {
  22. @Override
  23. public WaterSensor map(String value) throws Exception {
  24. String[] line = value.split(",");
  25. return new WaterSensor(
  26. line[0],
  27. Long.parseLong(line[1]),
  28. Integer.parseInt(line[2])
  29. );
  30. }
  31. })
  32. .assignTimestampsAndWatermarks(
  33. WatermarkStrategy
  34. .<WaterSensor>forMonotonousTimestamps()
  35. .withTimestampAssigner((data,ts)-> data.getTs() * 1000L)
  36. );
  37. KeyedStream<WaterSensor, String> sensorKS = sensorDS.keyBy(sensor -> sensor.getId());
  38. SingleOutputStreamOperator<String> resultDS = sensorKS.process(new KeyedProcessFunction<String, WaterSensor, String>() {
  39. @Override
  40. public void processElement(WaterSensor waterSensor, Context context, Collector<String> collector) throws Exception {
  41. TimerService timerService = context.timerService();
  42. System.out.println("currentProcessingTime=" + timerService.currentProcessingTime() + "\n"
  43. + "currentWatermark=" + timerService.currentWatermark());
  44. timerService.registerProcessingTimeTimer(System.currentTimeMillis() + 1000L);
  45. }
  46. /**
  47. * 定时器触发,当时间进展到了注册的时间,调用该方法
  48. * @param timestamp
  49. * @param ctx
  50. * @param out
  51. * @throws Exception
  52. */
  53. @Override
  54. public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {
  55. System.out.println("onTimer的timestamp=" + timestamp);
  56. }
  57. });
  58. resultDS.print();
  59. env.execute();
  60. }
  61. }

2.基于事件时间的定时器

  1. package com.atguigu.flink.day05;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.streaming.api.TimerService;
  6. import org.apache.flink.streaming.api.datastream.KeyedStream;
  7. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  8. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  9. import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
  10. import org.apache.flink.util.Collector;
  11. /**
  12. * 基于处理时间的定时器
  13. */
  14. public class $12_TimeServiceDemo2 {
  15. public static void main(String[] args) throws Exception {
  16. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  17. env.setParallelism(1);
  18. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  19. .map(new MapFunction<String, WaterSensor>() {
  20. @Override
  21. public WaterSensor map(String value) throws Exception {
  22. String[] line = value.split(",");
  23. return new WaterSensor(
  24. line[0],
  25. Long.parseLong(line[1]),
  26. Integer.parseInt(line[2])
  27. );
  28. }
  29. })
  30. .assignTimestampsAndWatermarks(
  31. WatermarkStrategy
  32. .<WaterSensor>forMonotonousTimestamps()
  33. .withTimestampAssigner((data,ts)-> data.getTs() * 1000L)
  34. );
  35. KeyedStream<WaterSensor, String> sensorKS = sensorDS.keyBy(sensor -> sensor.getId());
  36. SingleOutputStreamOperator<String> resultDS = sensorKS.process(new KeyedProcessFunction<String, WaterSensor, String>() {
  37. @Override
  38. public void processElement(WaterSensor waterSensor, Context context, Collector<String> collector) throws Exception {
  39. TimerService timerService = context.timerService();
  40. System.out.println("currentProcessingTime=" + timerService.currentProcessingTime() + "\n"
  41. + "currentWatermark=" + timerService.currentWatermark());
  42. timerService.registerEventTimeTimer(context.timestamp() + 5000L);
  43. }
  44. /**
  45. * 定时器触发,当时间进展到了注册的时间,调用该方法
  46. * @param timestamp
  47. * @param ctx
  48. * @param out
  49. * @throws Exception
  50. */
  51. @Override
  52. public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {
  53. System.out.println("onTimer的timestamp=" + timestamp +"\n"
  54. + "watermark=" + ctx.timerService().currentWatermark());
  55. }
  56. });
  57. resultDS.print();
  58. env.execute();
  59. }
  60. }

启动socket测试

  1. [atguigu@hadoop102 ~]$ nc -lk 9999
  2. s_1,1,1
  3. s_1,6,1
  4. s_1,7,1

day05[Flink流处理高阶编程(中)] - 图3