第一章.Flink的window机制

1.窗口概述及分类

  1. # 窗口概述
  2. 在流处理应用中,数据是连续不断的,因此我们不可能等到所有数据都到了才开始处理,当然我们可以每来一个消息就处理一次,但是有时我们需要做一些聚合类的运算,例如:在过去的1分钟内有多少用户点击了我们的网页,在这种情况下,我们必须定义一个窗口,用来收集最近一分钟内的数据,并对这个窗口中的数据进行计算
  3. 流式计算是一种被设计用于处理无限数据集的数据处理引擎,而无限数据集是指一种不断增长的本质上的无限的数据集,而window窗口是一种切割无限数据为有限块进行处理的手段
  4. 在Flink中,窗口(window)是处理无界流的核心,窗口把流切割成有限大小的多个存储桶(bucket),我们在这些桶上进行计算

窗口分为两类

  • 基于时间的窗口(时间驱动)
  • 基于元素个数的窗口(数据驱动)

2.基于时间的窗口

  • 时间窗口包含一个开始时间戳(包含)和结束时间戳(不包含)—->[start,end),这两个时间戳一起限制了窗口的尺寸,在Flink中,TimeWindow这个类表示基于时间的窗口
  • 时间窗口又分三种:滚动窗口,滑动窗口,会话窗口

一.滚动窗口(Tumbling Windows)

  • 滚动窗口有固定的大小,窗口与窗口之间不会重叠也没有缝隙,比如,如果指定一个长度为5分钟的窗口,当前窗口开始计算,每5分钟启动一个新的窗口
  • 滚动窗口能将数据流切分为不重叠的窗口,每一个事件只能属于一个窗口

day04[Flink流处理高阶编程(上)] - 图1

示例代码

  1. package com.atguigu.flink.day04;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.TimeCharacteristic;
  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.TumblingProcessingTimeWindows;
  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 org.junit.After;
  14. import org.junit.Before;
  15. import org.junit.Test;
  16. public class $01_TimeWindow {
  17. SingleOutputStreamOperator<WaterSensor> sensorDS = null;
  18. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = null;
  19. StreamExecutionEnvironment env = null;
  20. SingleOutputStreamOperator<String> resultDS = null;
  21. @Before
  22. public void before(){
  23. env = StreamExecutionEnvironment.getExecutionEnvironment();
  24. env.setParallelism(1);
  25. //老版本需要该设置
  26. //env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime);
  27. sensorDS = env.socketTextStream("hadoop102", 9999)
  28. .map(new MapFunction<String, WaterSensor>() {
  29. @Override
  30. public WaterSensor map(String value) throws Exception {
  31. String[] line = value.split(",");
  32. return new WaterSensor(
  33. line[0],
  34. Long.parseLong(line[1]),
  35. Integer.parseInt(line[2])
  36. );
  37. }
  38. });
  39. }
  40. /**
  41. * 基于时间的滚动窗口
  42. */
  43. @Test
  44. public void test01(){
  45. sensorWS = sensorDS.keyBy(sensor -> sensor.getId())
  46. /**
  47. * flink version <= 1.11的写法
  48. */
  49. //.timeWindow(Time.seconds(3));
  50. /**
  51. * 新写法
  52. */
  53. .window(TumblingProcessingTimeWindows.of(Time.seconds(3)));
  54. }
  55. @After
  56. public void after() throws Exception {
  57. 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. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  63. }
  64. });
  65. resultDS.print();
  66. env.execute();
  67. }
  68. }

二.滑动窗口(Sliding Windows)

  • 与滚动窗口一样,滑动窗口也有固定的长度,另外一个参数我们叫滑动步长,用来控制窗口启动的频率
  • 所以如果滑动步长小于窗口长度,滑动窗口会重叠,这种情况下,一个元素可能会被分配到多个窗口中

day04[Flink流处理高阶编程(上)] - 图2

示例代码

  1. /**
  2. * 基于时间的滑动窗口
  3. */
  4. @Test
  5. public void test02(){
  6. sensorWS = sensorDS.keyBy(sensor -> sensor.getId())
  7. /**
  8. * <=1.11的写法
  9. */
  10. //.timeWindow(Time.seconds(5),Time.seconds(3));
  11. /**
  12. * 新写法
  13. */
  14. .window(TumblingProcessingTimeWindows.of(Time.seconds(3)));
  15. }

三.会话窗口(Session Windows)

  • 会话窗口分配器会根据活动的元素进行分组,会话窗口不会有重叠,与滚动窗口和滑动窗口相比,会话窗口也没有固定的开启和关闭时间
  • 如果会话窗口有一段时间没有收到数据,会话窗口会自动关闭,这段没有收到数据的时间就是会话窗口的gap(间隔)
  • 我们可以配置静态的gap,也可以通过一个gap extractor 函数来定义gap的长度,当时间超过了这个gap,当前的会话窗口就会关闭,后续的元素会被分配到一个新的会话窗口

day04[Flink流处理高阶编程(上)] - 图3

示例代码

  1. /**
  2. * 基于时间的会话窗口-静态Gap
  3. */
  4. @Test
  5. public void test03(){
  6. sensorWS = sensorDS.keyBy(sensor -> sensor.getId())
  7. .window(ProcessingTimeSessionWindows.withGap(Time.seconds(5)));
  8. }
  9. /**
  10. * 基于时间的会话窗口-动态Gap
  11. */
  12. @Test
  13. public void test04(){
  14. sensorWS = sensorDS.keyBy(sensor -> sensor.getId())
  15. .window(ProcessingTimeSessionWindows.withDynamicGap(new SessionWindowTimeGapExtractor<WaterSensor>() {
  16. @Override
  17. public long extract(WaterSensor element) {
  18. return element.getTs() * 1000L;
  19. }
  20. }));
  21. }

3.全局窗口(Global Windows)

全局窗口分配器会分配相同key的所有元素,进入同一个Global window,这种窗口机制只有指定自定义的触发器时才有用,否则,不会做任何运算,因为这种窗口没有能够处理聚集在一起元素的结束点

day04[Flink流处理高阶编程(上)] - 图4

示例代码

  1. package com.atguigu.flink.day04;
  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.GlobalWindows;
  11. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  12. import org.apache.flink.streaming.api.windowing.evictors.TimeEvictor;
  13. import org.apache.flink.streaming.api.windowing.time.Time;
  14. import org.apache.flink.streaming.api.windowing.triggers.CountTrigger;
  15. import org.apache.flink.streaming.api.windowing.windows.GlobalWindow;
  16. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  17. import org.apache.flink.util.Collector;
  18. import java.util.Arrays;
  19. public class $02_GlobalWindow {
  20. public static void main(String[] args) throws Exception {
  21. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  22. env.setParallelism(1);
  23. /**
  24. * 老版本需要这么指定,默认是处理时间
  25. * 新版本不需要指定,默认是事件时间,而且开窗的时候,自己就要明确指定时间语义
  26. */
  27. //env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
  28. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  29. .map(new MapFunction<String, WaterSensor>() {
  30. @Override
  31. public WaterSensor map(String value) throws Exception {
  32. String[] line = value.split(",");
  33. return new WaterSensor(
  34. line[0],
  35. Long.parseLong(line[1]),
  36. Integer.parseInt(line[2])
  37. );
  38. }
  39. })
  40. .assignTimestampsAndWatermarks(
  41. WatermarkStrategy
  42. .<WaterSensor>forMonotonousTimestamps()
  43. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  44. @Override
  45. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  46. return element.getTs() * 1000L;
  47. }
  48. })
  49. );
  50. //TODO 全局窗口,需要自定义触发器
  51. WindowedStream<WaterSensor, String, GlobalWindow> sensorWS = sensorDS.
  52. keyBy(sensor -> sensor.getId())
  53. .window(GlobalWindows.create())
  54. .trigger(CountTrigger.of(3))
  55. .evictor(TimeEvictor.of(Time.seconds(5)));
  56. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, GlobalWindow>() {
  57. @Override
  58. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  59. out.collect("key=" + s + "\n"
  60. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  61. + "窗口划分:[" + context.window().toString()+ ")" + "\n\n"
  62. );
  63. }
  64. });
  65. resultDS.print();
  66. env.execute();
  67. }
  68. }

4.基于元素个数的窗口

一.滚动窗口

  • 默认的CountWindow是一个滚动窗口,只需要指定窗口大小即可,当元素数量达到窗口大小时,就会触发窗口的执行
  • 哪个窗口先达到3个元素,哪个窗口就关闭,不影响其他的窗口

示例代码

  1. package com.atguigu.flink.day04;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.datastream.WindowedStream;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  8. import org.apache.flink.streaming.api.windowing.windows.GlobalWindow;
  9. import org.apache.flink.util.Collector;
  10. public class $03_CountWindow {
  11. public static void main(String[] args) throws Exception {
  12. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  13. env.setParallelism(1);
  14. /**
  15. * 老版本需要这么指定,默认是处理时间
  16. * 新版本不需要指定,默认是事件时间,而且开窗的时候,自己就要明确指定时间语义
  17. */
  18. //env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
  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. //TODO 基于条数的滚动窗口
  32. WindowedStream<WaterSensor, String, GlobalWindow> sensorWS = sensorDS.
  33. keyBy(sensor -> sensor.getId())
  34. .countWindow(3);
  35. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, GlobalWindow>() {
  36. @Override
  37. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  38. out.collect("key=" + s + "\n"
  39. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  40. + "窗口划分:[" + context.window().toString()+ ")" + "\n\n"
  41. );
  42. }
  43. });
  44. resultDS.print();
  45. env.execute();
  46. }
  47. }

二.滑动窗口

滑动窗口和滚动窗口的函数名是完全一致的,只是在传参数时需要传入两个参数,一个是window_size,一个是sliding_size

示例代码

  1. package com.atguigu.flink.day04;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.datastream.WindowedStream;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  8. import org.apache.flink.streaming.api.windowing.windows.GlobalWindow;
  9. import org.apache.flink.util.Collector;
  10. /**
  11. * 窗口的划分并不是以第一条数据为基准的
  12. * 第一个输出的窗口 = 步长
  13. * 每经过一个步长,都有一个窗口输出,关闭
  14. */
  15. public class $04_CountWindow2 {
  16. public static void main(String[] args) throws Exception {
  17. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.setParallelism(1);
  19. /**
  20. * 老版本需要这么指定,默认是处理时间
  21. * 新版本不需要指定,默认是事件时间,而且开窗的时候,自己就要明确指定时间语义
  22. */
  23. //env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
  24. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  25. .map(new MapFunction<String, WaterSensor>() {
  26. @Override
  27. public WaterSensor map(String value) throws Exception {
  28. String[] line = value.split(",");
  29. return new WaterSensor(
  30. line[0],
  31. Long.parseLong(line[1]),
  32. Integer.parseInt(line[2])
  33. );
  34. }
  35. });
  36. //TODO 基于条数的滑动窗口
  37. WindowedStream<WaterSensor, String, GlobalWindow> sensorWS = sensorDS.
  38. keyBy(sensor -> sensor.getId())
  39. .countWindow(5,2);
  40. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, GlobalWindow>() {
  41. @Override
  42. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  43. out.collect("key=" + s + "\n"
  44. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  45. + "窗口划分:[" + context.window().toString()+ ")" + "\n\n"
  46. );
  47. }
  48. });
  49. resultDS.print();
  50. env.execute();
  51. }
  52. }

5.窗口函数

  • 指定了窗口的分配器后,需要指定如何计算,这事由窗口函数来负责,一旦窗口关闭,窗口函数去计算处理窗口中的每个元素
  • 窗口函数可以是ReduceFunction,AggregateFunction,ProcessWindowFunction中的任意一种
  • ReduceFunction,AggregateFunction更加高效,原因是Flink可以对到来的元素进行增量聚合,ProcessWindowFunction可以得到一个包含这个窗口的所有元素的迭代器,以及这些元素所属窗口的元数据信息
  • ProcessWindowFunction不能被高效执行的原因是Flink在执行这个函数之前,需要在内存缓存这个窗口上的所有元素

一.增量函数

增量函数:

  • 来一条处理一条
  • 窗口的输出次数:窗口触发关闭的时候,也就是说,每个窗口只会向下游传递一次结果
  • aggregate相比reduce:都是增量聚合
  • reduce第一条不调用方法
  • aggregate比reduce更灵活,reduce要求输入类型,中间状态,输出类型要一致,aggregate输入类型,中间状态(累加器),输出类型可以不一致

示例代码

  1. package com.atguigu.flink.day04;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.AggregateFunction;
  4. import org.apache.flink.api.common.functions.MapFunction;
  5. import org.apache.flink.api.common.functions.ReduceFunction;
  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.windowing.assigners.TumblingProcessingTimeWindows;
  10. import org.apache.flink.streaming.api.windowing.time.Time;
  11. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  12. import org.junit.After;
  13. import org.junit.Before;
  14. import org.junit.Test;
  15. public class $05_IncreFunction {
  16. SingleOutputStreamOperator<WaterSensor> sensorDS = null;
  17. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = null;
  18. StreamExecutionEnvironment env = null;
  19. SingleOutputStreamOperator<WaterSensor> resultDS = null;
  20. SingleOutputStreamOperator<String> resultDS1 = null;
  21. @Before
  22. public void before(){
  23. env = StreamExecutionEnvironment.getExecutionEnvironment();
  24. env.setParallelism(1);
  25. sensorDS = env.socketTextStream("hadoop102", 9999)
  26. .map(new MapFunction<String, WaterSensor>() {
  27. @Override
  28. public WaterSensor map(String value) throws Exception {
  29. String[] line = value.split(",");
  30. return new WaterSensor(
  31. line[0],
  32. Long.parseLong(line[1]),
  33. Integer.parseInt(line[2])
  34. );
  35. }
  36. });
  37. sensorWS = sensorDS.keyBy(sensor -> sensor.getId())
  38. .window(TumblingProcessingTimeWindows.of(Time.seconds(10)));//窗口分配器
  39. }
  40. /**
  41. * 窗口函数-增量函数 - reduce
  42. */
  43. @Test
  44. public void test01(){
  45. resultDS = sensorWS.reduce(new ReduceFunction<WaterSensor>() {
  46. @Override
  47. public WaterSensor reduce(WaterSensor value1, WaterSensor value2) throws Exception {
  48. System.out.println(value1 + "<=======>" + value2);
  49. return new WaterSensor(value1.getId(), System.currentTimeMillis(), value1.getVc()+value2.getVc());
  50. }
  51. });
  52. }
  53. /**
  54. * 窗口函数-增量函数 -aggregate
  55. */
  56. @Test
  57. public void test02(){
  58. resultDS1 = sensorWS.aggregate(new AggregateFunction<WaterSensor, Integer, String>() {
  59. /**
  60. * 初始化累加器
  61. * @return
  62. */
  63. @Override
  64. public Integer createAccumulator() {
  65. System.out.println("create.....");
  66. return 0;
  67. }
  68. /**
  69. * 累加的逻辑
  70. * @param value
  71. * @param accumulator
  72. * @return
  73. */
  74. @Override
  75. public Integer add(WaterSensor value, Integer accumulator) {
  76. System.out.println("add....");
  77. return value.getVc() + accumulator;
  78. }
  79. /**
  80. * 获取结果
  81. * @param accumulator
  82. * @return
  83. */
  84. @Override
  85. public String getResult(Integer accumulator) {
  86. System.out.println("getResult....");
  87. return accumulator.toString();
  88. }
  89. /**
  90. * 只有会话窗口才会调用
  91. * @param a
  92. * @param b
  93. * @return
  94. */
  95. @Override
  96. public Integer merge(Integer a, Integer b) {
  97. System.out.println("merge.....");
  98. return a + b;
  99. }
  100. });
  101. }
  102. @After
  103. public void after() throws Exception {
  104. if(resultDS != null){
  105. resultDS.print();
  106. }
  107. if(resultDS1 != null){
  108. resultDS1.print();
  109. }
  110. env.execute();
  111. }
  112. }

二.全量函数

全量函数:数据全部存起来,触发的时候一次性计算和输出

示例代码

  1. package com.atguigu.flink.day04;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.datastream.WindowedStream;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  8. import org.apache.flink.streaming.api.windowing.assigners.TumblingProcessingTimeWindows;
  9. import org.apache.flink.streaming.api.windowing.time.Time;
  10. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  11. import org.apache.flink.util.Collector;
  12. import java.util.Arrays;
  13. /**
  14. * 窗口函数-->全量函数
  15. */
  16. public class $06_AllFunctionWindow {
  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. //TODO 全局窗口,需要自定义触发器
  33. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  34. keyBy(sensor -> sensor.getId())
  35. .window(TumblingProcessingTimeWindows.of(Time.seconds(10)));
  36. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  37. @Override
  38. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  39. out.collect("key=" + s + "\n"
  40. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  41. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  42. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n"
  43. );
  44. }
  45. });
  46. resultDS.print();
  47. env.execute();
  48. }
  49. }

6.Non-Keyed Windows

  • 在用窗口前首先需要确认应该是在keyBy之后的流上用,还是在没有keyBy的流上用
  • 在keyed streams上使用窗口,窗口计算被并行的运用在多个task上,可以认为每个task都有自己单独窗口
  • 在非non-keyed stream上使用窗口,流的并行度只能是1,所有的窗口逻辑只能在一个单独的task上执行
  • 需要注意的是:非key分区的流上使用window,如果把并行度强行设置为>1,则会抛出异常

示例代码

  1. package com.atguigu.flink.day04;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.functions.MapFunction;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.streaming.api.functions.windowing.ProcessAllWindowFunction;
  7. import org.apache.flink.streaming.api.windowing.assigners.TumblingProcessingTimeWindows;
  8. import org.apache.flink.streaming.api.windowing.time.Time;
  9. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  10. import org.apache.flink.util.Collector;
  11. import java.util.Arrays;
  12. /**
  13. * NoKeyedWindow
  14. */
  15. public class $07_NoKeyedWindow {
  16. public static void main(String[] args) throws Exception {
  17. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.setParallelism(3);
  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. sensorDS.windowAll(TumblingProcessingTimeWindows.of(Time.seconds(10)))
  32. .process(new ProcessAllWindowFunction<WaterSensor, String, TimeWindow>() {
  33. @Override
  34. public void process(Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  35. out.collect(Arrays.asList(elements).toString());
  36. }
  37. }).print();
  38. env.execute();
  39. }
  40. }

7.窗口源码分析

以滚动窗口为例

  1. package com.atguigu.flink.day04;
  2. /**
  3. * 窗口相关源码分析(以滚动窗口为例)
  4. */
  5. public class $09_SourceCodeAnalysis {
  6. public static void main(String[] args) throws Exception {
  7. /**
  8. *
  9. * 窗口是怎么划分的
  10. * long start =TimeWindow.getWindowStartWithOffset(
  11. * now, (globalOffset + staggerOffset) % size, size);
  12. * return Collections.singletonList(new TimeWindow(start, start + size));
  13. * public static long getWindowStartWithOffset(long timestamp, long offset, long windowSize) {
  14. * return timestamp - (timestamp - offset + windowSize) % windowSize;
  15. * }
  16. *
  17. */
  18. /**
  19. *
  20. * 窗口为什么左闭右开
  21. * Gets the largest timestamp that still belongs to this window.
  22. *
  23. * <p>This timestamp is identical to {@code getEnd() - 1}.
  24. *
  25. * @return The largest timestamp that still belongs to this window.
  26. * @see #getEnd()
  27. *
  28. *
  29. * @Override
  30. * public long maxTimestamp() {
  31. * return end - 1;
  32. * }
  33. */
  34. /**
  35. * 窗口什么时候触发
  36. * @Override
  37. * public TriggerResult onProcessingTime(long time, TimeWindow window, TriggerContext ctx) {
  38. * return TriggerResult.FIRE;
  39. * }
  40. * 时间 >= maxTimestamp
  41. */
  42. /**
  43. * 窗口什么时候创建?窗口什么时候销毁?(窗口的生命周期?)
  44. * 创建: 属于本窗口的第一条数据来的时候,new的,放到一个 单例集合
  45. * 销毁: cleanupTime = window.maxTimestamp() + allowedLateness;
  46. * 清空状态的时间 = 时间进展到 窗口最大时间戳 + 允许迟到的时间
  47. */
  48. }
  49. }

第二章.Flink的时间语义与WaterMark

1.时间语义

一.事件时间

  1. 事件时间是指这个事件发生的时间
  2. 在event进入flink之前,通常被嵌入到了event中,一般作为这个event的时间戳存在
  3. 在事件时间体系中,时间的进度依赖于事件本身,和任何设备的时间无关,事件时间必须制定如何产生Event Time Watermarks(水印),在事件时间体系中,水印是代表时间进度的标志(作用相当于现实时间的时钟)
  4. 在理想条件下,不管事件时间何时到达或者他们到达的顺序如何,事件时间处理将产生完全一致且确定的结果,事件时间会在等待无序时间(迟到时间)时产生一定的延迟,由于只能等待有限的时间,因此这限制了确定性事件时间应用程序的可使用性
  5. 假设所有数据都已到达,事件时间操作将按照预期方式进行,即使在处理无序或迟到的事件或重新处理历史数据时,也会产生正确且一致的效果,例如,每小时事件时间窗口将包含带有事件时间戳的所有记录,将记录计入该小时,无论他们到达的顺序或处理时间
  6. 在使用窗口的时候,如果使用事件时间,就指定时间分配器为事件时间分配器

注意: 在1.12之前默认的时间语义是处理时间,从1.12开始,Flink内部已经把默认的语义改成了事件时间

二.处理时间

  1. 处理时间指的是执行操作的各个设备的时间
  2. 对于运行在处理时间上的流程序,所有的基于时间的操作(比如时间窗口)都是使用的设备时钟,比如一个长度为1个小时的窗口将会包含设备时钟表示的1个小时内所有的数据,假设应用程序在9:15am启动,第一个小时窗口将会包含9:15am到10:00am所有的数据,然后下个窗口是10:00am-11.00am等
  3. 处理时间是最简单时间语义,数据流与设备之间不需要做任何的协调,它提供了最好的性能和最低的延迟,但是在分布式和异步的环境之下,处理时间没有方法保证确定性,容易受到数据传递速度的影响:事件的延迟和乱序
  4. 在使用窗口的时候,如果使用处理时间,就指定时间分配器为处理时间分配器

三.案例

day04[Flink流处理高阶编程(上)] - 图5

四.代码

  1. package com.atguigu.flink.day04;
  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.TimeCharacteristic;
  7. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  8. import org.apache.flink.streaming.api.datastream.WindowedStream;
  9. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  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.assigners.TumblingProcessingTimeWindows;
  13. import org.apache.flink.streaming.api.windowing.assigners.TumblingTimeWindows;
  14. import org.apache.flink.streaming.api.windowing.time.Time;
  15. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  16. import org.apache.flink.util.Collector;
  17. import java.util.Arrays;
  18. public class $08_TimeCharacteristic_EventTime {
  19. public static void main(String[] args) throws Exception {
  20. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  21. env.setParallelism(1);
  22. /**
  23. * 老版本需要这么指定,默认是处理时间
  24. * 新版本不需要指定,默认是事件时间,而且开窗的时候,自己就要明确指定时间语义
  25. */
  26. //env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
  27. SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop102", 9999)
  28. .map(new MapFunction<String, WaterSensor>() {
  29. @Override
  30. public WaterSensor map(String value) throws Exception {
  31. String[] line = value.split(",");
  32. return new WaterSensor(
  33. line[0],
  34. Long.parseLong(line[1]),
  35. Integer.parseInt(line[2])
  36. );
  37. }
  38. })
  39. .assignTimestampsAndWatermarks(
  40. WatermarkStrategy
  41. .<WaterSensor>forMonotonousTimestamps()
  42. .withTimestampAssigner(new SerializableTimestampAssigner<WaterSensor>() {
  43. @Override
  44. public long extractTimestamp(WaterSensor element, long recordTimestamp) {
  45. return element.getTs() * 1000L;
  46. }
  47. })
  48. );
  49. WindowedStream<WaterSensor, String, TimeWindow> sensorWS = sensorDS.
  50. keyBy(sensor -> sensor.getId())
  51. .window(TumblingEventTimeWindows.of(Time.seconds(10)));
  52. SingleOutputStreamOperator<String> resultDS = sensorWS.process(new ProcessWindowFunction<WaterSensor, String, String, TimeWindow>() {
  53. @Override
  54. public void process(String s, Context context, Iterable<WaterSensor> elements, Collector<String> out) throws Exception {
  55. out.collect("key" + s + "\n"
  56. + "窗口有" + elements.spliterator().estimateSize() + "条数据" + "\n"
  57. + "数据为:" + Arrays.asList(elements).toString() + "\n"
  58. + "窗口划分:[" + context.window().getStart() + "," + context.window().getEnd() + ")" + "\n\n");
  59. }
  60. });
  61. resultDS.print();
  62. env.execute();
  63. }
  64. }

2.WaterMark初理解

  • WaterMark是用来衡量事件时间的进展
  • 用来解决乱序的问题
  • 是一个特殊的时间戳,从指定生成的位置插入到流里
  • 单调不减的
  • 用来触发窗口的
  • Flink认为,事件时间小于WaterMark的数据应该都已经处理过了,如果后续还有事件时间小于watermark的数据来,称为迟到数据
  1. 支持event time的流式处理框架需要一种能够测量 event time 进度的方式,比如一个窗口算子创建了一个长度为1小时的窗口,那么这个算子需要知道事件时间已经到达了这个窗口的关闭时间,从而在程序中去关闭这个窗口
  2. 事件时间可以不依赖处理时间来表示时间的进度,例如在程序中,即使处理时间和事件时间有相同的速度,事件时间可能会轻微的落后处理时间,另外一方面,使用事件时间可以在几秒内处理已经缓存在kafka中的数据,这些数据可以照样被正确处理,就像实时发生的一样能够进入正确的窗口
  3. 这种在flink中去测量事件时间的进度的机制就是waterMark(水印),waterMark作为数据流的一部分在流动,并且携带一个时间戳t
  4. 一个WaterMark(t)表示在这个流里面事件时间已经到了时间t,意味着此时,流中不应该存在这样的数据:它的时间戳t2<=t(时间比较旧或者等于时间戳)

WaterMark原理图

day04[Flink流处理高阶编程(上)] - 图6

day04[Flink流处理高阶编程(上)] - 图7