You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

关于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.21 07:04:30