算子实例的事件时间是否可能变小?基于Flink并行流水印的技术问询
关于Flink窗口事件时间与水印顺序的问题解答
首先,先明确Flink并行流水印的核心规则:
算子的当前事件时间为其所有输入流事件时间的最小值
接下来咱们拆解你提到的两个场景:
场景1:水印29先到达,之后水印14到达
这里要先纠正一个小误解:每个输入流的水印是单调递增的,而且算子的事件时间是取所有输入流当前水印的最小值。
当水印29先到达窗口(1)时,另一个输入流(也就是会发水印14的那个流)此时还没发送任何水印,或者它的水印还停留在比14更小的初始值(比如负无穷)。这时候窗口的事件时间会取min(29, 初始值),也就是初始值,而不是直接设为29。
当之后水印14到达时,这个输入流的水印更新为14(注意:这个流的水印必须是单调递增的,所以14肯定比它之前的水印大),此时窗口的事件时间就变成min(29,14)=14——这看起来像是事件时间“变小”,但本质是因为我们之前只拿到了一个流的水印,现在两个流的水印都齐了,取到了真正的最小值而已。
场景2:之后数据源(2)生成的水印39到达
当数据源(2)的水印39到达时,这个流的水印更新为39(同样,39必须比它之前的14大,符合单调递增规则)。此时窗口的两个输入流水印分别是29(来自第一个流)和39(来自第二个流),所以窗口的事件时间会取min(29,39)=29。只有当第一个流的水印也更新到大于等于29(比如到30)时,窗口的事件时间才会跟着上升到30。
总结一下:窗口的事件时间始终是所有输入流当前水印的最小值,而且单个输入流的水印只会涨不会跌,所以窗口的事件时间整体趋势是单调递增的——偶尔看起来“下降”,只是因为之前没拿到所有输入流的水印,现在补全后取到了更小的有效值而已。
内容的提问来源于stack exchange,提问作者YuFeng Shen
相关产品推荐
相关产品推荐

