Flink旧Checkpoint无法正确删除的问题排查求助
问题原因分析
1. Checkpoint超时导致的未完成文件残留
当Checkpoint因超时失败时,Flink的Checkpoint协调器无法完成完整的清理流程:RocksDB生成Checkpoint(尤其是增量Checkpoint)时会先写入大量临时文件,这些文件仅在Checkpoint成功完成后才会被_metadata文件标记为有效。一旦Checkpoint超时中断,未被_metadata关联的临时文件会被遗留在S3的chk目录中,成为无主垃圾文件。
2. S3最终一致性的影响
S3属于最终一致性存储,Flink在判断文件是否被_metadata引用时,可能因S3的读写延迟出现判断偏差:例如Flink读取_metadata时,文件内容尚未完全同步,导致部分已写入的文件未被识别为待清理的垃圾。
3. Flink 1.15.1版本的已知缺陷
Flink 1.15.x系列在S3作为RocksDB后端的场景下,存在Checkpoint清理逻辑的bug:当Checkpoint失败后,垃圾文件回收机制未被正确触发,尤其是在作业持续运行且多次出现Checkpoint超时的情况下,残留文件会持续累积。
解决方法
1. 优化Checkpoint配置,降低失败概率
- 调整Checkpoint超时时间:增大
execution.checkpointing.timeout参数,避免因作业负载过高频繁触发Checkpoint超时。 - 配置外部化Checkpoint的清理策略:设置
execution.checkpointing.externalized-checkpoint-retention为DELETE_ON_CANCELLATION,确保作业取消时自动清理无效Checkpoint文件;若需保留失败Checkpoint用于排查,可设为RETAIN_ON_CANCELLATION。 - 合理设置保留的Checkpoint数量:通过
state.checkpoints.num-retained控制有效Checkpoint的保留数量,Flink会自动清理超出数量的旧Checkpoint。
2. 触发主动垃圾回收
- 使用Flink内置工具:运行
CheckpointFilesCleaner工具,指定S3上的Checkpoint根目录,工具会自动解析每个_metadata文件,清理未被引用的垃圾文件。执行命令示例:./bin/flink run -c org.apache.flink.runtime.state.filesystem.CheckpointFilesCleaner \ -Dstate.checkpoints.dir=s3://your-checkpoint-bucket/path \ -Dstate.backend=rocksdb - 自定义脚本清理:编写脚本遍历S3上的所有
chk-<id>目录,读取每个目录下的_metadata文件,提取其中引用的文件列表,删除目录中不在列表内的文件。操作前需备份,避免误删有效文件。
3. 升级Flink版本
升级至Flink 1.16及以上稳定版本,后续版本修复了多个S3 Checkpoint清理相关的bug,包括增量Checkpoint的垃圾文件回收逻辑、S3路径处理的一致性问题等,从根源上减少残留文件的产生。
4. 安全手动清理
若需手动删除垃圾文件,可先暂停作业(或停止数据源),确认当前有效的Checkpoint目录(即Flink UI中显示的最新成功Checkpoint),删除其他未被引用的文件。删除后重启作业,验证作业能正常从有效Checkpoint恢复即可。
内容的提问来源于stack exchange,提问作者C.S.
相关产品推荐
相关产品推荐

