Flink使用Join算子时如何正确配置并推进Watermark?
疑问解答
- 是,Flink中每个算子的并行实例都会持有独立的Watermark生成器,单条流的全局有效Watermark为该流所有并行实例上报Watermark的最小值,只要有一个实例没有收到事件推进本地Watermark,整个流的Watermark就不会向前推进。
- 不需要。Watermark跟随流数据在算子之间传播,和keyBy的先后顺序没有直接关系,同Key的事件后续会被路由到同一个Join算子实例,但Join算子收到的两个流的Watermark,还是分别取各自流所有上游并行实例的最小Watermark,只有两个流的Watermark都超过窗口结束时间才会触发窗口计算。
异常现象原因
你当前问题的核心是assignTimestampsAndWatermarks算子并行度为2时,前4条测试事件(0s的A、B事件,62s的A、B事件)都被路由到了其中一个并行实例,另一个并行实例始终没有收到任何事件,它的本地Watermark一直停留在初始值Long.MIN_VALUE,导致两个流的全局Watermark都被卡住无法超过60s,1分钟滚动窗口自然不会触发计算。
直到发送122s的事件时,事件被分配到了之前空闲的第二个并行实例,该实例的本地Watermark推进到121s,此时两个流的全局Watermark才超过60s,之前的0s事件所在的窗口才触发Join输出。
修复方案
1. 调整Watermark赋值位置
如果你的事件是按Key分区输入的,可以把assignTimestampsAndWatermarks放到keyBy之后,这样每个Key对应的事件都会进入同一个Watermark生成器实例,避免出现空闲实例:
// 修改后的流处理逻辑 sourceA = initialSourceA.map(parseToEvent).keyBy(Event.Key) streamA = sourceA.assignTimestampsAndWatermarks(CustomWatermarkStrategy()) sourceB = initialSourceB.map(parseToEvent).keyBy(Event.Key) streamB = sourceB.assignTimestampsAndWatermarks(CustomWatermarkStrategy())
2. 增加空闲并行实例的Watermark推进逻辑
如果你的Source本身就会出现部分并行实例长期没有数据的情况,需要在Watermark策略中增加空闲超时配置,当实例空闲超过指定时间后,自动把该实例标记为空闲,不会再卡住全局Watermark:
// 赋值Watermark时增加空闲超时配置 .assignTimestampsAndWatermarks( WatermarkStrategy .forGenerator(CustomWatermarkStrategy()) .withIdleness(Duration.ofSeconds(5)) // 空闲5s就标记该实例为空闲 )
3. 优化自定义Watermark生成器逻辑
你当前的生成器每个事件都触发一次Watermark上报,会产生大量不必要的Watermark传播开销,建议改为在onPeriodicEmit方法中定期上报Watermark,onEvent只更新当前最大时间戳即可:
class CustomWaterMarkGenerator( private val maxOutOfOrderness: Long, private var currentMaxTimeStamp: Long = Long.MIN_VALUE, ) : WatermarkGenerator<EventType> { override fun onEvent(event: EventType, eventTimestamp: Long, output: WatermarkOutput) { currentMaxTimeStamp = currentMaxTimeStamp.coerceAtLeast(eventTimestamp) } override fun onPeriodicEmit(output: WatermarkOutput) { output.emitWatermark(Watermark(currentMaxTimeStamp - maxOutOfOrderness - 1)); } }
Flink默认的周期性上报间隔是200ms,足够满足绝大多数场景的实时性要求,同时大幅降低Watermark传播开销。
Join场景Watermark使用注意事项
- 两个流的
assignTimestampsAndWatermarks需要独立配置,确保两个流的Watermark都能正常推进 - 窗口Join的触发条件是两个流的全局Watermark都超过窗口的结束时间,任意一个流的Watermark卡住都会导致窗口无法触发
- 多并行度场景下一定要处理空闲并行实例的Watermark推进问题,避免全局Watermark被空闲实例卡住
内容的提问来源于stack exchange,提问作者Akula

