本文主要是介绍flink watermark处理细节-StatusWatermarkValve代码分析,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!
首先抛出一个问题:
kafka topic下有3个partition,下游consumer为flink job,flink job的并行度为4,如下图
那么window operator的watermark是否会一直很小,导致窗口迟迟不触发计算
理清这个问题需要看flink对watermark的处理,StatusWatermarkValve类嵌入了Watermark和StreamStatus两种元素怎么发送到下游到逻辑,inputStreamStatus方法包含了主要的处理逻辑
public void inputStreamStatus(StreamStatus streamStatus, int channelIndex)
这篇关于flink watermark处理细节-StatusWatermarkValve代码分析的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!