flink watermark处理细节-StatusWatermarkValve代码分析

it2025-03-20  15

首先抛出一个问题:

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)
最新回复(0)