无流入流量时Flink滚动窗口Checkpoint大小持续增长问题咨询
环境与配置
Flink版本:1.20
Checkpoint与状态后端配置
# Checkpoint settings checkpoint.timeout=600000 checkpoint.interval=5000 checkpoint.pause-between=10000 # State backend settings state.backend.type=rocksdb state.backend.incremental=true state.backend.rocksdb.checkpoint.transfer.thread.num=8 state.backend.rocksdb.thread.num=8 state.backend.rocksdb.use-bloom-filter=true state.checkpoint-storage=filesystem
基础工作流代码
KafkaSource<MyEvent> kafkaSource = KafkaSource.<MyEvent>builder() .setBootstrapServers(brokerList) .setTopics(sourceTopic) .setGroupId(groupId) .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.LATEST)) .setValueOnlyDeserializer(new JsonDeserializationSchema<>(MyEvent.class)) .build(); SingleOutputStreamOperator<MyEvent> eventStream = env .fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source") .uid("Kafka Source"); SingleOutputStreamOperator<ProcessedResult> processedStream = eventStream .setMaxParallelism(4) .filter(new EventFilterFunction()) // Custom filtering logic .uid("Event Filter Function") .windowAll(TumblingProcessingTimeWindows.of(Duration.ofSeconds(5))) // Windowing logic .trigger(new CustomTimerTrigger<>()) // Custom trigger .process(new GroupAndAggregateWindowFunction()) // Aggregation logic .uid("Group And Aggregate Window Function") .name("Group And Aggregate Window Function"); processedStream .process(new ComputationFunction()) // Further processing .sinkTo(kafkaSink); // Write results to Kafka
窗口处理函数代码
public class GroupAndAggregateWindowFunction extends ProcessAllWindowFunction<MyEvent, AggregatedResult, TimeWindow> { @Override public void process(Context context, Iterable<MyEvent> events, Collector<AggregatedResult> out) throws Exception { Map<CompositeKey, List<MyEvent>> groupedEvents = new HashMap<>(); for (MyEvent event : events) { CompositeKey key = new CompositeKey(event.getField1(), event.getField2()); groupedEvents.computeIfAbsent(key, k -> new ArrayList<>()).add(event); } for (Map.Entry<CompositeKey, List<MyEvent>> entry : groupedEvents.entrySet()) { CompositeKey key = entry.getKey(); List<MyEvent> eventGroup = entry.getValue(); if (eventGroup.isEmpty()) continue; AggregatedResult result = new AggregatedResult(); result.setField1(key.getField1()); result.setField2(key.getField2()); List<Item> aggregatedItems = new ArrayList<>(); for (MyEvent event : eventGroup) { if (event.getItems() != null) { aggregatedItems.addAll(event.getItems()); } } if (aggregatedItems.isEmpty()) continue; result.setItems(aggregatedItems); out.collect(result); } } }
自定义触发器代码
public class CustomTimerTrigger<T> extends Trigger<T, TimeWindow> { @Override public TriggerResult onElement(Object element, long timestamp, TimeWindow window, TriggerContext ctx) throws Exception { ctx.registerProcessingTimeTimer(window.getEnd()); return TriggerResult.CONTINUE; // Don't fire yet } @Override public TriggerResult onProcessingTime(long time, TimeWindow window, TriggerContext ctx) { return TriggerResult.FIRE_AND_PURGE; // Fire and clear the window } @Override public TriggerResult onEventTime(long time, TimeWindow window, TriggerContext ctx) { return TriggerResult.CONTINUE; // No event-time processing } @Override public void clear(TimeWindow window, TriggerContext ctx) { ctx.deleteProcessingTimeTimer(window.getEnd()); // Clean up } }
问题描述
使用5秒滚动窗口结合ProcessAllWindowFunction做聚合,发现即使函数里没有显式声明状态,聚合算子的Checkpoint大小仍每日增长约10MB,无流入流量时也是如此。
原本以为添加了返回FIRE_AND_PURGE的自定义触发器后,无流量时Checkpoint大小会完全下降,但实际仅下降0.1-0.5MB,且持续维持在10MB以上并逐日累积。
提出以下问题:
- 为何
ProcessAllWindowFunction无显式状态,Checkpoint大小仍持续增长? - 如何确保无新记录流入时Flink正确清理窗口状态?
- RocksDB配置或Checkpoint配置是否会影响状态保留?
问题解答
1. 为何ProcessAllWindowFunction无显式状态,Checkpoint大小仍持续增长?
ProcessAllWindowFunction无需显式定义状态,但Flink的窗口机制会自动维护窗口内的元素状态——所有进入窗口的事件都会被Flink存储到状态中,直到窗口被清理。
另外还有几个关键原因:
- 自定义触发器的
FIRE_AND_PURGE虽会触发计算并清理窗口状态,但如果窗口触发不及时(比如处理时间漂移)、触发器元数据(如注册的ProcessingTimeTimer)未彻底清理,会导致状态残留; - RocksDB增量Checkpoint会保留历史版本数据,即使状态被清理,旧的Checkpoint文件不会立即删除,直到Checkpoint保留策略触发清理;
windowAll是全局窗口,所有事件进入同一个窗口实例,若清理逻辑存在遗漏,状态会持续累积。
2. 如何确保无新记录流入时Flink正确清理窗口状态?
可以从以下几个方向优化:
- 验证触发器执行逻辑:在触发器的
onProcessingTime和clear方法中添加日志,确认窗口结束时这两个方法都被正常执行,确保定时器被删除、窗口状态被清理; - 显式设置窗口允许延迟:为滚动窗口设置
allowedLateness(Duration.ZERO),确保窗口结束后不再接收元素,也不会保留状态; - 配置状态TTL:为窗口状态设置过期时间,即使清理逻辑有遗漏,TTL到期后状态会被自动清理,示例配置:
StateTtlConfig ttlConfig = StateTtlConfig.newBuilder(Duration.ofMinutes(1)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); env.getConfig().setGlobalStateTtlConfig(ttlConfig); - 无流量时触发全量Checkpoint:手动触发一次全量Checkpoint,增量Checkpoint可能需要全量Checkpoint才能体现状态的实际减少;
- 检查算子并行度配置:
windowAll并行度为1,确保setMaxParallelism设置合理,避免状态分片残留。
3. RocksDB配置或Checkpoint配置是否会影响状态保留?
是的,部分配置会直接影响状态保留:
- 增量Checkpoint与保留策略:
state.backend.incremental=true会让RocksDB只存储与上一个Checkpoint的差异,若state.checkpoints.num-retained设置过大,旧Checkpoint文件会持续占用空间,看起来像是状态增长; - RocksDB压缩配置:未开启压缩(如
state.backend.rocksdb.compression.type=lz4)会导致磁盘上的SST文件体积增大; - Checkpoint频率:
checkpoint.interval=5000(5秒一次)频率过高,增量Checkpoint会生成大量小文件,累积后占用空间; - Checkpoint超时:
checkpoint.timeout=600000(10分钟)若超时失败,可能导致部分状态未被正确清理。
建议调整:
- 将
state.checkpoints.num-retained设置为合理值(如3),避免保留过多旧Checkpoint; - 开启RocksDB压缩,减少磁盘占用;
- 无流量场景下降低Checkpoint频率,或临时停止Checkpoint。
内容的提问来源于stack exchange,提问作者Kendo
相关产品推荐
相关产品推荐

