Flink Time
Event Time与Process Time
如果你对数据的准确性要求比较高的话,采用 Event time 能保障 exactly-once。Processing Time 一般用于实时消费、精准性要求略低的场景,主要是因为时间生成不是 deterministic。
我们可以看下面的关系图, X 轴是 Event time,Y 轴是 Processing time。理想情况下 Event time 和 Processing Time 是相同的,就是说只要有一个事件发生,就可以立刻处理。但是实际场景中,事件发生后往往会经过一定延时才会被处理,这样就会导致我们系统的时间往往会滞后于事件时间。这里它们两个的差 Processing-time lag 表示我们处理事件的延时。

事件时间常用在窗口中,使用 watermark 来确保数据完备性,比如说 watermarker 值大于 window 末尾时间时,我们就可以认为 window 窗口所有数据都已经到达了,就可以触发计算了。

比如上面 [0-10] 的窗口,现在 watermark 走到了 10,已经到达了窗口的结束,触发计算 SUM=21。如果要是想对迟到的数据再进行触发,可以再定义一下后面 late data 的触发,比如说后面来了个 9,我们的 SUM 就等于 30。
北京时间与格林尼治时间

划定时间窗口和时间戳以毫秒为单位
1 | eventStream.assignTimestampsAndWatermarks(assigner = new AscendingTimestampExtractor[Event]() { |
processFunction中的时间
1 | new ProcessFunction[Event, (String, Long)] { |
1 | context.timestamp=1565064599999 |
以事件事件,窗口大小5min
一条消息,事件的时间是 1565064350(2019-08-06 12:05:50)Context.timestamp是1565064599999(2019-08-06 12:09:59), 可以理解为窗口结束的时间context.timerService().currentWatermark() 莫名其妙,这么大context.timerService().currentProcessingTime()是当前处理事件 2019-08-06 12:10:07
Context.timestamp
1 | Timestamp of the element currently being processed or timestamp of a firing timer. |
时间戳timestamp与水印watermark
提取Timestamp与生成Watermark一般步骤
- 设置时间特性为Event Time。
StreamExecutionEnvironment#setStreamTimeCharacteristic(TimeCharacteristic.EventTime)。 - 在Source后Window前用
DataStream#assignTimestampsAndWatermarks方法(AssignerWithPeriodicWatermarks或AssignerWithPunctuatedWatermarks)提取时间戳并生成水印。 - 重写
extractTimestamp方法提取Timestamp,重写getCurrentWatermark方法或checkAndGetNextWatermark方法生成水印。
数据源中指定
1 | public static class ExampleSourceFunction implements SourceFunction<Tuple4<String,Long,String,Integer>>{ |
时间戳和水位线分别通过两个方法实现的
时间戳的抽取方法
1 | /** |
周期性水印
分配器为元素分配事件时间时间戳,并生成表示流中事件时间进度的低水位线。这些时间戳和水印由对事件时间进行操作的函数和运算符使用,例如事件时间窗口。
使用此类以周期性间隔生成水印。最多每隔I毫秒(通过ExecutionConfig.getautowerMarkinterval())配置一次,系统将调用getCurrentWatermark()方法来探测下一个水印值。
如果探测到的值为非空,并且时间戳大于前一个水印的时间戳,系统将生成新的水印(以保留升序水印的约定)。 如果自上次调用getCurrentWatermark()方法后没有新元素到达,系统调用该方法的频率可能会低于每I毫秒一次。
时间戳和水印被定义为代表自纪元(世界协调时1970年1月1日午夜)以来的毫秒的长度。某个值为t的水印表示不会再出现带有事件时间戳x的元素,其中x小于或等于t。
1 | /** |
非连续性水印
1 | /** |

升序水印 AscendingTimestampExtractor
1 | org.apache.flink.streaming.api.functions.timestamps.AscendingTimestampExtractor |
AscendingTimestampExtractor抽象类实现了ssignerWithPeriodicWatermarks接口的extractTimestamp及getCurrentWatermark方法,同时声明抽象方法extractAscendingTimestamp供子类实现
水印的生产策略为:getCurrentWatermark 方法在currentTimestamp不为Long.MIN_VALUE时返回Watermark(currentTimestamp - 1)
综上:
发现AscendingTimestampExtractor适用于elements的时间在每个并行task里面事单调递增的,(timestamp monotony)数据场景。
无序水印 BoundedOutOfOrdernessTimestampExtractor
1 | public abstract class BoundedOutOfOrdernessTimestampExtractor<T> implements AssignerWithPeriodicWatermarks<T> { |
BoundedOutOfOrdernessTimestampExtractor 抽象类实现AssignerWithPeriodicWatermarks接口的extractTimestamp及getCurrentWatermark方法,同时声明抽象方法extractAscendingTimestamp供子类实现BoundedOutOfOrdernessTimestampExtractor 的构造器接收maxOutOfOrderness参数用于指定element允许滞后(t~`t_w,t为element的eventTime,t_w为前一次watermark的时间)的最大时间,在计算窗口数据时,如果超过该值则会被忽略。 BoundedOutOfOrdernessTimestampExtractor的extractTimestamp方法会调用子类的extractTimestamp方法抽取时间,如果该时间大于currentMaxTimestamp,则更新currentMaxTimestamp; getCurrentWatermark先计算potentialWM,如果potentialWM大于等于lastEmittedWatermark则更新lastEmittedWatemakr(currentMaxTimestamp - lastEmittedWatermark >= maxOutOfOrderness, 这里表示lastEmittedWatermark太小了,所以差值超过了maxOutOfOrderness,因此会调大lastEmittedWatermark),最后返回watermark`
具体逻辑参考:就是判断如果新的水印时间大于上次最后一次发出水印时间,选择新的水印时间,否则选择上次最后发送水印时间戳
TimeStamp分配器和Watermark生成器(Timestamp Assigners / Watermark Generators)
数据延迟
主要的办法是给定一个允许延迟的时间,在该时间范围内仍可以接受处理延迟数据
设置允许延迟的时间是通过allowedLateness(lateness: Time)设置
保存延迟数据则是通过sideOutputLateData(outputTag: OutputTag[T])保存
获取延迟数据是通过DataStream.getSideOutput(tag: OutputTag[X])获取
allowedLateness(lateness: Time)
1 | def allowedLateness(lateness: Time): WindowedStream[T, K, W] = { |
该方法传入一个Time值,设置允许数据迟到的时间,这个时间和waterMark中的时间概念不同。再来回顾一下,
waterMark=数据的事件时间-允许乱序时间值
随着新数据的到来,waterMark的值会更新为最新数据事件时间-允许乱序时间值,但是如果这时候来了一条历史数据,waterMark值则不会更新。总的来说,waterMark是为了能接收到尽可能多的乱序数据。
那这里的Time值呢?主要是为了等待迟到的数据,在一定时间范围内,如果属于该窗口的数据到来,仍会进行计算,后面会对计算方式仔细说明
allowedLateness 是等待数据进入窗口的最大时间,未进入的数据,可以通过偏流导出
注意:该方法只针对于基于event-time的窗口,如果是基于processing-time,并且指定了非零的time值则会抛出异常
sideOutputLateData(outputTag: OutputTag[T])
1 | def sideOutputLateData(outputTag: OutputTag[T]): WindowedStream[T, K, W] = { |
该方法是将迟来的数据保存至给定的outputTag参数,而OutputTag则是用来标记延迟数据的一个对象。
DataStream.getSideOutput(tag: OutputTag[X])
通过window等操作返回的DataStream调用该方法,传入标记延迟数据的对象来获取延迟的数据
1 | waterStream.keyBy(0) |
【参考文献】