Flink 1.17升级后Checkpoint大小持续增长的原因分析
问题描述
我们近期将若干Flink服务从1.13版本升级至1.17版本,未进行其他重大代码变更。应用程序每2分钟执行一次Checkpoint。升级后,发现Checkpoint大小持续增长,仅在不使用最近Checkpoint重新部署应用时才会重置。橙棕色线条为Flink 1.13版本期间的Checkpoint大小,之后为Flink 1.17版本,图中的下降点代表重新部署并从头启动的节点。所有状态均设置了TTL,所有窗口均配置了对应的聚合函数。
可能的原因
状态TTL清理逻辑变更导致过期状态堆积
Flink 1.13到1.17版本间,状态TTL的触发机制有调整。比如默认清理策略从ON_CREATE_AND_WRITE改为ON_READ_AND_WRITE后,若部分键控状态长期无读写操作,过期状态不会被主动清理,仅在下次访问时触发,导致大量过期状态持续留在Checkpoint中。另外,窗口状态的TTL清理需配合allowedLateness等参数生效,若升级后参数未适配新版本逻辑,窗口状态也无法及时释放。窗口状态管理逻辑调整
新版本优化了窗口的关闭、清理流程,比如窗口触发计算后的状态清理时机延迟,或需显式配置cleanupAfterTimeWindow等参数才能触发清理。即便配置了聚合函数,若未同步调整窗口相关参数,已完成计算的窗口状态可能持续占用Checkpoint空间。RocksDB状态后端行为变化
若使用RocksDB状态后端,1.17版本对其集成做了更新:默认压缩策略、sst文件合并逻辑调整,导致过期状态的物理文件未被及时合并删除;或状态清理触发频率降低,使得Checkpoint需包含更多未清理的历史数据,最终引发体积持续增长。增量Checkpoint适配问题
Flink 1.17优化了增量Checkpoint实现,但升级后若未正确配置incremental-checkpoint.enabled等参数,或新版本增量Checkpoint在特定场景下无法有效复用历史快照,会导致每次Checkpoint都需写入更多新数据,进而使体积持续上升。状态序列化与元数据变化
新版本更新了状态序列化框架或Checkpoint元数据格式,若旧状态无法被兼容清理,或元数据中额外存储的信息未及时清理历史条目,也会导致Checkpoint体积逐渐增大。
内容的提问来源于stack exchange,提问作者FelixNavidad

