关于Apache Flink中Watermark传播逻辑的疑问
关于Apache Flink中Watermark传播逻辑的疑问
嗨,我来帮你把Union算子的水印传播逻辑讲清楚,这确实是Flink事件时间处理里容易绕晕的点~
首先要记住一个核心规则:当多个流通过Union算子合并后,下游算子收到的水印始终是所有输入流当前水印中的最小值。这是Flink为了保证事件时间一致性的设计——只有当所有输入流都确认“某个时间点之前的事件已经全部到达”,全局的事件时间进度才会推进到这个点,避免因为某条流的水印超前,导致其他流的迟到数据被错误丢弃。
下面针对你的三个场景逐一分析(假设你用的是默认的单调递增水印生成器,即水印值等于流中已收到的最大事件时间,没有额外的延迟配置):
第一个场景
- Source One输出事件时间为10的消息,此时它的水印更新为10
- Source Two输出事件时间为9的消息,此时它的水印更新为9
- Union后的水印取两者最小值,也就是9,你的猜测是对的。
第二个场景
- Source One输出事件时间为12的消息,水印更新为12(水印是单调递增的,不会回退)
- Source Two输出事件时间为11的消息,水印更新为11
- Union后的水印取当前两个流水印的最小值,也就是11,而不是之前的9——因为每个流的水印只会随着新事件的到来向前推进,不会停留在旧值。
第三个场景
- Source One输出事件时间为15的消息,水印更新为15
- Source Two输出事件时间为13的消息,水印更新为13
- Union后的水印依然取最小值,也就是13。哪怕Source Two有延迟,只要它的水印还没追上Source One,Union后的整体水印就会被它拖慢,直到Source Two的水印继续推进。
简单总结一下:Union算子的水印是“取所有输入流当前水印的最小值”,而且每个流的水印是单调递增的,只会向前走不会回头,所以每次新事件到来后,先更新对应流的水印,再取最小值作为Union后的输出水印。
备注:内容来源于stack exchange,提问作者Fred
相关产品推荐
相关产品推荐

