关于Flink中BoundedOutOfOrder Watermarks的原理与迟滞数据问询
Flink 有界乱序水印(BoundedOutOfOrder Watermarks)工作机制确认
事件时序说明
11:00 11:01 11:02 11:03 11:04 11:05 11:06 11:07 11:08 <=== 处理时间(Processing Time) | | |--------1---------2----------3---------4--------5--------------------------------- <== 输入事件(INPUT EVENTS) | | (水印) |--------1------------2----------3-----------------5------W!-----------------4 <==== 事件时间(Event Time) 11:01 11:02 11:03 11:05 11:05 11:04
使用的水印策略配置
WatermarkStrategy .<Tuple2<Long, String>>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.f0);
问题1:Flink会每5秒向管道中注入一个Watermark事件?
回答:错误
forBoundedOutOfOrderness水印策略并非固定间隔注入水印,而是基于输入事件的事件时间戳动态生成:
- 水印的时间戳计算公式为:
当前观察到的最大事件时间戳 - 允许的乱序时长(此处为5秒) - 只有当新事件的事件时间戳推动了当前最大事件时间戳,且计算出的新水印值大于上一次发送的水印值时,才会向下游发送新水印
- 若没有新事件输入,不会生成任何水印
问题2:event-4在Watermark关闭窗口后到达,是否会被视为迟滞数据,除非配置allowedLateness否则将被丢弃?
回答:是的
结合你的场景:
- 目标窗口为
[11:00, 11:05)(左闭右开),当水印时间戳≥窗口结束时间(11:05)时,窗口会被触发计算并关闭 - event-4的事件时间为11:04,属于该窗口,但它在窗口关闭后才到达,属于迟到数据
- 默认情况下,未配置
allowedLateness时,这类迟到数据会直接被丢弃;若配置了允许的迟到时长(如allowedLateness(Duration.ofSeconds(2))),则在窗口关闭后的指定时长内到达的迟到数据,仍会被纳入窗口并触发窗口重新计算
内容的提问来源于stack exchange,提问作者Mandar K
相关产品推荐
相关产品推荐

