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

无流入流量时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以上并逐日累积。

提出以下问题:

  1. 为何ProcessAllWindowFunction无显式状态,Checkpoint大小仍持续增长?
  2. 如何确保无新记录流入时Flink正确清理窗口状态?
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 04:20:02