未使用状态时TumblingProcessingTimeWindows导致Checkpoint持续增大问题问询
问题分析与解决:Flink Checkpoint持续增大+TaskManager线程增长
核心现象
- 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 []
问题根源
- 窗口数据未被自动清理:使用
ProcessingTime窗口且未设置watermark时,Flink的窗口清理机制依赖水印推进判断窗口是否过期。由于没有水印,Flink无法确认窗口是否彻底结束,所有窗口的输入数据会一直保留在状态中,直接导致Checkpoint持续增大。 - 线程增长的关联逻辑:未清理的窗口持续积累,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
相关产品推荐
相关产品推荐

