Flink Checkpoint 越来越大怎么办?从 200MB 涨到 12GB 的完整排查记录(附 State TTL 配置)

🔑 关键词:Flink Checkpoint过大, Flink状态膨胀排查, RocksDB StateBackend调优, State TTL配置, Flink大状态

📖 摘要:一个日均 800 万单的 Flink 去重作业,Checkpoint 从上线时的 200MB 一路涨到 12GB,耗时从 3 秒变成 92 秒。本文记录完整排查过程、三个关键命令、State TTL 具体参数,以及为什么我劝你先别动 RocksDB 配置。

Flink Checkpoint 越来越大怎么办?从 200MB 涨到 12GB 的完整排查记录

图片

先说结论,免得你跟我一样走弯路:Checkpoint 体积膨胀,绝大多数情况不是 RocksDB 参数没调好,是你在 KeyedState 里存了不该存的东西。我这次遇到的是一个日均 800 万单的去重逻辑,ValueState<Boolean> 存了 30 天从来不删,状态从上线时的 200MB 一路涨到 12GB,Checkpoint 耗时从 3 秒变成 92 秒,作业平均每 48 小时挂一次,从 Savepoint 恢复要 7 分 40 秒。

先把环境交代清楚,不然对比没意义:Flink 1.17.1,YARN Application 模式;3 个 TaskManager,每个 8 核 16GB,managed memory 按默认占 40%;状态后端是 EmbeddedRocksDBStateBackend,开了增量 Checkpoint,state.backend.incremental: true;Checkpoint interval 60s,timeout 10min,exactly-once,Kafka source 并行度 36;数据侧峰值 12 万条/秒,日均 800 万笔订单事件,key 用的是 order_id,字符串长度 32 字节左右。

一、我一开始走的两条弯路,你可以直接跳过

第一条弯路是调 RocksDB。我把 state.backend.rocksdb.block.cache-size 从 64MB 调到 256MB,把 writebuffer.size 从 64MB 调到 128MB,重启,盯着监控看了两小时。Checkpoint 大小一字节没变。这不是废话吗——block cache 影响的是读性能和读放大,跟状态体积没有半毛钱关系。我当时就是看见「RocksDB」四个字就条件反射了,运维同学在群里 @ 我说第三次 checkpoint 超时的时候,我脑子是空的。

第二条弯路是换状态后端。有同事说 HashMapStateBackend 快,快照也快。我试了,作业起来 40 秒后直接被 YARN 干掉——12GB 的状态要塞进 JVM 堆,而 TM 总共才 16GB,还要留给网络缓冲、算子对象、元空间、Direct Memory。这个方案在小状态场景下确实没毛病,状态上了 10GB 就是自寻死路,OOMKilled 是必然的。顺便说一句,很多人以为换状态后端能解决 Checkpoint 慢,其实只解决快照的序列化方式,解决不了状态本身有多大。

图片

二、真正定位问题的三个命令

第一个是去看 Checkpoint 历史里的两个数。 Flink Web UI → Job → Checkpoints → History,每次记录里有两个容易看混的指标:State Size(这是全量状态大小)和 Checkpointed Data Size(这是本次快照实际上传的增量数据)。开了增量 Checkpoint 之后这俩能差 10 倍以上。我当时 Checkpointed Data Size 一直是 380MB 左右,看着挺健康,但 State Size 已经 12.1GB 了。只盯增量数据的人,会一直以为自己的作业没问题。

第二个是上机器看 RocksDB 的磁盘占用。

du -sh /data/flink/io/tmp/*/rocksdb/*

图片

注意这个路径跟着 io.tmp.dirs 走,不同集群不一样。看到单 TM 占用 41GB 的时候我就确定,问题在状态本身,不在快照流程。

第三个是把 RocksDB 的原生指标打开。

state.backend.rocksdb.metrics.num-keys: true
state.backend.rocksdb.metrics.estimate-live-data-size: true

estimate-live-data-size 这个指标特别有用,它直接告诉你 RocksDB 自己认为有多少活数据。我这边跑出来 10.8GB,跟 Checkpoint 的 12.1GB 基本对得上,差距是索引和元数据的开销。

三、根因:一个从来没被清理过的 ValueState

图片

翻代码翻了半小时,找到了。业务方要求「同一个订单只处理一次」,实现方式是每条数据来的时候查一下 ValueState<Boolean>,没有就置 true 并往下走。逻辑没错,问题是这个 state 从来没设过过期时间。上线第一天当然只有当天数据,跑了 30 天之后,里面躺着 2.4 亿个 key。

这里有个估算公式,我在后面又验证过两次,误差在 20% 以内,够用了:

状态大小 ≈ Σ(key 平均字节 + value 序列化字节 + RocksDB 固定开销约 60~80 字节) × (1 − LZ4 压缩率)

代进去算:32 字节 key + 1 字节 value + 75 字节 SST 内部开销 ≈ 108 字节/条;2.4 亿条 × 108 字节 × 0.5 压缩率 ≈ 12.9GB。跟观测到的 12.1GB 吻合。

修复配置长这样:

图片

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.days(3))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .cleanupInRocksdbCompactFilter(10000L)
    .build();

ValueStateDescriptor<Boolean> desc =
    new ValueStateDescriptor<>("orderDedup", Boolean.class);
desc.enableTimeToLive(ttlConfig);

cleanupInRocksdbCompactFilter 的参数是「每处理多少个状态条目触发一次时间戳过滤」,默认 1000。设太小会明显拖慢 compaction,设太大清理不及时,我这边 10000 是试出来的折中值。还有个坑必须提醒:TTL 只在读写和 compaction 时生效,如果某个 key 再也不来了,它不会被主动删除。所以 du -sh 下降是滞后的,可能滞后好几个小时,别盯着看,会焦虑。

四、修完之后的对比,以及三个我觉得反直觉的结论

图片

指标 修复前 修复后
State Size 12.1 GB 1.7 GB
Checkpointed Data Size(增量) 380 MB 45 MB
Checkpoint 耗时 P99 92 s 8 s
从 Savepoint 恢复 7 min 40 s 55 s
单 TM 磁盘占用 41 GB 6 GB
稳定运行时长 约 48 h 挂一次 连续 60 天以上

第一个结论:Checkpoint 大小本质是「业务记忆长度 × 数据密度」的投影,它是架构指标,不是运维指标。运维怎么调参数都调不动它,因为它记录的是你业务上决定要记多久。

第二个结论:TTL 设多久不是技术问题,是业务问题。我一开始自己拍了 7 天,结果业务方说对账周期是 T+3,那 7 天就是纯浪费。这种事应该让产品经理回答,而不是我们拍脑袋。

第三个结论:增量 Checkpoint 会掩盖问题。它让你每次看增量数据都很小,很容易误判。建议把 State Size 单独接进 Prometheus 做个增长率告警,比如 7 天涨了 30% 就报警,而不是等作业挂了才查。

那什么时候不该做 TTL?如果状态确实是「每一条都不能丢」,比如金融的幂等记账,那就别硬上 TTL。可以把状态外置到 Redis 或 HBase,但代价要算清楚:Flink 内置状态是微秒级访问,Redis 大概 0.5~2ms,HBase 5~20ms。去重、累计这类「每条数据都要查一次」的场景,外置方案延迟直接翻 100 倍,别轻易用。只有维表关联这种低频查询的才适合外置。最后一句大实话:先去翻代码看有没有没设 TTL 的 state,再回来谈调参。

🏷️ 标签: