0%

Flink watermark

1 watermark

a window operator registers a timer for every active window, which cleans up the window’s state when the event time passes the window’s ending time.

1.1 处理空闲数据源

在某些情况下,由于数据产生的比较少,导致一段时间内没有数据产生,进而就没有水印的生成,导致下游依赖水印的一些操作就会出现问题,
比如某一个算子的上游有多个算子,这种情况下,水印是取其上游两个算子的较小值,如果上游某一个算子因为缺少数据迟迟没有生成水印,就会出现eventtime倾斜问题,导致下游没法触发计算。

所以filnk通过WatermarkStrategy.withIdleness()方法允许用户在配置的时间内(即超时时间内)没有记录到达时将一个流标记为空闲。这样就意味着下游的数据不需要等待水印的到来。

当下次有水印生成并发射到下游的时候,这个数据流重新变成活跃状态。

通过下面的代码来实现对于空闲数据流的处理

1
2
3
WatermarkStrategy
.<Tuple2<Long, String>>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withIdleness(Duration.ofMinutes(1));

在 Flink SQL 使用 Kafka 作为 Source 的场景中,即使已配置 40 秒乱序 Watermark,仍可能出现 Watermark 滞后数分钟甚至 10 分钟以上 的现象;这在 多并行度单并行度 情况下本质原因一致,均不属于乱序参数失效。Watermark 的推进规则是:算子 Watermark 由其输入 Channel 的最小值决定,且只能在消费到新数据时基于最新事件时间推进,不会随系统时间自动前进。多并行度时,低频或空闲的 Kafka partition 会拖慢整体 Watermark(需通过 scan.idle-timeout 处理 idle 分区);单并行度时,则更应重点关注 Kafka 是否存在历史或批量数据、当前是否持续有新数据到达、event_time 字段本身是否整体晚于当前时间。此外,任务在 failover / savepoint 恢复 后会继承历史 Watermark 状态,短时间内表现为明显滞后;下游算子或 Sink 出现 反压、阻塞、慢 IO 或重 UDF 计算,也可能导致 Source 无法持续 poll 数据,从而间接延缓 Watermark 发送。需要明确的是,乱序时间只用于容忍事件顺序的乱序,并不等价于限制 Watermark 最大延迟;当 Watermark 长时间落后时,应系统性地从 Kafka 分区与入流模式、事件时间字段的真实性与时效性、任务恢复历史以及上下游执行链路是否阻塞等方面进行排查并沉淀为工程经验。