Flink延迟15分钟输出的事件为何能被下游Sink接收?
针对你遇到的Flink作业中延迟事件仍被Sink接收的情况,核心原因可以从以下几点拆解:
Timer触发后的输出不受当前Watermark的限制
你通过ProcessFunction注册的是事件时间Timer,当对应Key的Watermark推进到Timer设定的触发时间(事件时间+15分钟)时,Timer会触发并执行你定义的输出逻辑。这个输出动作是Timer触发后的主动行为,和当前Sink或全局的Watermark是否高于该事件的原始事件时间无关——Watermark的作用是触发Timer、窗口这类基于事件时间的计算,而非阻断已经触发的输出操作。Sink默认不做事件时间过滤
Flink的Sink组件默认只负责接收上游算子发送的数据,不会根据自身的Watermark值去判断是否丢弃事件时间小于当前Watermark的数据。只有上游的算子(比如窗口算子结合allowedLateness配置)会处理迟到数据,而你的场景中是ProcessFunction主动输出Timer触发的数据,Sink自然会照单全收。Watermark的分区特性可能存在认知偏差
注意Flink中Watermark是按Key分区推进的(KeyBy之后),全局或Sink看到的最低Watermark是所有分区Watermark的最小值。但你触发Timer的那个Key对应的分区,其Watermark是刚好推进到Timer触发时间才会执行输出,此时该分区的Watermark并不高于Timer触发的事件时间(即原始事件时间+15分钟),而Sink的最低Watermark可能是其他分区的低Watermark拉低的,但这不影响当前Key分区的Timer输出。
内容的提问来源于stack exchange,提问作者user12331

