You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

未使用状态时TumblingProcessingTimeWindows导致Checkpoint持续增大问题问询

核心现象

  • Checkpoint大小持续上升无回落,且大小与TumblingProcessingTimeWindows的Bytes Received几乎持平
  • TaskManager线程数持续增长
  • 业务逻辑正常执行,但未在WindowFunction中显式使用状态存储

关联代码片段

ret_stream = (log_stream
              .map(MyMapFunction())
              .filter(lambda x: self.get_key(x) is not None)
              .key_by(self.get_key, key_type=Types.STRING())
              .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
              .apply(MyWindowFunction()))

class MyWindowFunction(WindowFunction[tuple, tuple, str, TimeWindow]:
    def apply(self, key: KEY, window: W, inputs: Iterable[IN]) -> Iterable[OUT]:
        logger.info({
            "key": key,
            "data": [data for _, _, data in inputs]
        })
        return []

问题根源

  1. 窗口数据未被自动清理:使用ProcessingTime窗口且未设置watermark时,Flink的窗口清理机制依赖水印推进判断窗口是否过期。由于没有水印,Flink无法确认窗口是否彻底结束,所有窗口的输入数据会一直保留在状态中,直接导致Checkpoint持续增大。
  2. 线程增长的关联逻辑:未清理的窗口持续积累,Flink需要维护越来越多的窗口状态元数据和数据,后台窗口清理、状态维护线程负载上升,同时Checkpoint过程中会生成更多临时线程,最终引发线程数持续增长。

解决方案

方案1:为窗口设置明确的迟到数据规则

显式配置allowedLateness,让Flink在窗口触发后主动清理过期窗口:

ret_stream = (log_stream
              .map(MyMapFunction())
              .filter(lambda x: self.get_key(x) is not None)
              .key_by(self.get_key, key_type=Types.STRING())
              .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
              .allowed_lateness(Time.seconds(0))  # 不允许迟到数据,窗口触发后立即清理
              .apply(MyWindowFunction()))

若需处理少量迟到数据,可设置合理延迟时长(如5秒),Flink会在延迟期结束后清理窗口状态。

方案2:改用ProcessWindowFunction主动清理状态

WindowFunction底层仍依赖Flink状态存储保存窗口数据,改用ProcessWindowFunction可主动控制状态清理:

class MyProcessWindowFunction(ProcessWindowFunction[tuple, tuple, str, TimeWindow]):
    def process(self, key: str, context: Context, elements: Iterable[tuple]) -> Iterable[tuple]:
        logger.info({
            "key": key,
            "data": [data for _, _, data in elements]
        })
        # 主动清理当前窗口的状态
        context.window_state().clear()
        return []

ret_stream = (log_stream
              .map(MyMapFunction())
              .filter(lambda x: self.get_key(x) is not None)
              .key_by(self.get_key, key_type=Types.STRING())
              .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
              .process(MyProcessWindowFunction()))

方案3:启用ProcessingTime水印生成器

通过自动生成ProcessingTime水印,让Flink感知时间推进以触发窗口清理:

from flink.common.watermark_strategy import ProcessingTimeWatermarkStrategy

log_stream = log_stream.assign_timestamps_and_watermarks(
    ProcessingTimeWatermarkStrategy.for_monotonous_timestamps()
)

验证建议

  • 修改配置后观察Checkpoint大小,正常情况下会在窗口清理后回落至合理范围
  • 监控TaskManager线程数,确认不再持续增长

内容的提问来源于stack exchange,提问作者jimmy_go

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 13:06:20