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

Flink有状态作业Checkpoint持续膨胀问题咨询

我运行着一个Flink有状态作业,其Checkpoint大小随时间持续增长,最终导致Checkpoint操作出现问题。该状态采用ListStateDescriptor类型,代码片段如下:

public class MessageMapper extends
        KeyedProcessFunction<String, Message, Message> {
    private transient ListState<Trip> trips;

    @Override
    public void open(Configuration parameters) {
        ListStateDescriptor<Trip> tripDescriptor =
                new ListStateDescriptor<>("trips", Trip.class);
        trips = getRuntimeContext().getListState(tripDescriptor);
        ...
    }

    @Override
    public void processElement(Message message, Context ctx,
                               Collector<Message> out) throws Exception {

        List<Trip> tripList = new ArrayList<>();
        trips.get().forEach(tripList::add);


        // 根据消息内容,对状态执行以下操作之一:
        // - 向tripList添加新条目
        // - 从tripList移除现有条目

        // 更新状态
        trips.update(tripList);
    }
}

处理新消息时,trips状态平均应包含3-5个条目,条目会持续新增和删除,我预期状态大小保持稳定,但疑惑为何Checkpoint大小在数周/数月内持续增长,且数据分区的唯一Key数量未发生变化。

现咨询两个问题:

  1. 执行List.update()时,旧列表会被从存储中清除,还是会永久保留(进而导致Checkpoint大小增长)?
  2. 即便我不在意数据是否保留数月,是否仍需配置TTL?

问题解答

1. List.update()对旧列表的处理

执行trips.update(tripList)时,旧的列表状态会被替换,不会永久保留在当前活跃的状态中。但Checkpoint持续增长的核心原因大概率是代码存在状态泄露:

  • 逻辑上的条目删除操作可能存在疏漏,比如判断条件错误、对象匹配失败,导致实际没有移除目标Trip,列表条目悄悄累积;
  • 若使用全量Checkpoint模式,每次Checkpoint都会写入当前完整状态,只要实际状态的真实大小(比如条目数远超预期的3-5个)在缓慢增加,Checkpoint体积就会持续膨胀。

另外,Flink默认会保留最近若干个历史Checkpoint版本,但这只会导致有限的体积增长,不会出现数月持续变大的情况,所以优先排查代码逻辑中的状态累积问题。

2. 是否需要配置TTL

即便不在意数据保留时长,仍然建议配置状态TTL,原因如下:

  • 自动清理因代码疏漏产生的“僵尸状态”,比如某些Key不再有数据流入,但状态仍留在系统中;
  • 兜底防范意外场景(如某个Key的条目未被正确删除、持续累积)导致的状态无限膨胀;
  • 配合状态后端的清理机制,自动回收无效状态空间,从根源上避免Checkpoint持续增长的风险。

配置TTL的示例代码如下:

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.days(7))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .build();

ListStateDescriptor<Trip> tripDescriptor = new ListStateDescriptor<>("trips", Trip.class);
tripDescriptor.enableTimeToLive(ttlConfig);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 00:22:27