- Published on
Flink Checkpoint 生命周期管理:从 HDFS 清理到 DolphinScheduler 资源一致性
- Authors

- Name
- Charles Chen
在 Flink 生产环境中,Checkpoint 不只是一个定时生成的目录,而是作业故障恢复链路的一部分。它既要在故障时可靠可用,又不能无限占用 HDFS。
问题通常出现在作业取消、重新发布或调整启动脚本之后:旧作业实例留下的 Checkpoint 目录不会自动消失,新作业实例又会创建新的目录。运行次数越多,历史状态越多,最终形成一批已经没有恢复价值、却仍占用大量 HDFS 空间的数据。
本文以一次实际处理过程为例,说明:
- 为什么配置
RETAIN_ON_CANCELLATION后会留下历史 Checkpoint; - 如何用 Python 和 DolphinScheduler 管理旧目录;
- 为什么 HDFS 已经删除,DolphinScheduler 资源中心仍可能显示文件;
- 大型生产环境通常如何设计 Checkpoint 生命周期,而不是只依赖“按时间删除”。
一、为什么需要管理 Checkpoint 生命周期
假设 Checkpoint 根目录最初设计为:
hdfs:///flink/checkpoints/
├── flinkjob1/
│ ├── e4563582bfd39f062d20fe34d3093797/
│ ├── 195b5bf921b2526ede6e9e884807ad57/
│ └── 107f6daa627fc555e17903894ecd2c69/
└── flinkjob2/
├── 1/
├── 2/
└── 3/
第一层是逻辑作业名,第二层代表不同运行批次或 Job ID。作业取消后再次提交,会产生新的运行目录。
当配置如下策略时:
ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION
Flink 在作业取消时会保留外部化 Checkpoint,使后续仍能从它恢复。这正是该配置的价值,但代价是这些状态需要由用户自行管理和清理。Flink 官方 API 也明确指出,该模式会保留 Checkpoint 的元数据和状态,清理责任由用户承担。
另一个容易混淆的配置是:
execution.checkpointing.num-retained: 1
它控制的是一个正在运行的作业保留多少个已完成 Checkpoint,默认值为 1。它不是整个 HDFS 根目录的全局保留上限,也不会自动清除此前已取消作业遗留的外部化 Checkpoint。
因此,生命周期管理要解决的不是“Flink 为什么没有遵守保留数量”,而是“谁来管理已经脱离当前运行作业控制的历史恢复点”。
二、Checkpoint、Externalized Checkpoint 和 Savepoint 的边界
生产治理前必须先分清三种对象:
| 对象 | 主要用途 | 生命周期 |
|---|---|---|
| 普通 Checkpoint | 当前作业自动故障恢复 | 通常由 Flink 管理 |
| Externalized Checkpoint | 取消或失败后仍允许从状态恢复 | 取决于保留策略,可能需要人工清理 |
| Savepoint | 计划发布、升级、迁移和人工回滚 | 明确由用户管理 |
Checkpoint 更偏向自动故障恢复,Savepoint 更适合作为计划变更的发布工件。生产环境不应把“永远保留所有 Checkpoint”当成发布回滚体系。
一个更清晰的做法是:
- 作业运行期间保留少量最近 Checkpoint;
- 发布或重大变更前主动创建 Savepoint;
- 为发布版本记录可恢复路径和保留截止时间;
- 超过回滚窗口后,再清理旧运行批次。
三、清理脚本的安全边界
这次采用 Python 3 脚本执行 HDFS 操作,因为相较于纯 Shell,它更容易处理目录排序、参数校验、保护列表、预演模式和错误退出。脚本自身只依赖 Python 标准库,通过 hdfs dfs 命令访问 HDFS。
清理规则是:
- 遍历 Checkpoint 根目录下的每个逻辑作业;
- 读取每个作业下的直属子目录;
- 按 HDFS 修改时间排序;
- 始终保留最新的一个运行目录;
- 其余目录只有超过最小年龄后才成为候选项;
- 默认只打印,只有显式传入
--execute才执行删除。
配套脚本位于 flink_checkpoint_lifecycle.py。
为什么不能直接删除单个 chk-*
Flink 状态目录可能包含 shared/、taskowned/、元数据以及增量状态文件。不同 Checkpoint 之间可能共享状态文件。绕过 Flink 的引用关系,直接挑选单个 chk-* 删除,存在破坏剩余恢复点的风险。
本文脚本的清理单位是已经退出的旧作业实例目录,不是当前作业内的单个 Checkpoint。即使如此,生产执行前仍必须确认该目录不属于活跃作业,也没有被发布回滚流程保护。
四、在 DolphinScheduler 中运行
先把 Python 脚本上传到 DolphinScheduler 资源中心,并在 Shell 节点中绑定该资源。任务启动时,Worker 会将绑定资源下载到本次任务的临时执行目录,因此 Shell 节点需要查找的是本地临时目录中的 Python 文件,不是 HDFS 上的 Checkpoint 根目录。
下面是可直接放入 DolphinScheduler Shell 节点的预演任务:
#!/bin/bash
set -eu
SCRIPT_PATH="$(find "$PWD" -name 'flink_checkpoint_lifecycle.py' -type f -print -quit)"
CKPT_ROOT="hdfs:///flink/checkpoints/test"
if [ -z "$SCRIPT_PATH" ]; then
echo "ERROR: cleanup script was not downloaded"
exit 1
fi
echo "script=$SCRIPT_PATH"
echo "checkpoint_root=$CKPT_ROOT"
python3 --version
python3 "$SCRIPT_PATH" \
--root "$CKPT_ROOT" \
--min-age-hours 72 \
--max-list-lines 1000
这里:
set -e:任一命令失败后立即让任务失败,避免错误被吞掉;set -u:引用未定义变量时立即失败;SCRIPT_PATH:查找 DolphinScheduler 下载到 Worker 本地的资源文件;CKPT_ROOT:真正需要治理的 HDFS Checkpoint 根目录;--min-age-hours 72:仅把至少 72 小时前的旧目录列为候选;--max-list-lines 1000:限制单次 HDFS 列表输出规模;- 未传
--execute:只预演,不删除。
预演日志确认无误后,生产删除只增加一个参数:
python3 "$SCRIPT_PATH" \
--root "$CKPT_ROOT" \
--min-age-hours 72 \
--max-list-lines 1000 \
--execute
日志中的关键状态含义如下:
| 输出 | 含义 |
|---|---|
KEEP | 当前作业下最新的目录,始终保留 |
SKIP-YOUNG | 不是最新目录,但年龄还没有超过阈值 |
WOULD-DELETE | 预演模式下计划删除 |
DELETE | 执行模式下正在删除 |
PROTECT | 被保护参数明确排除的路径 |
定时任务不建议一开始就每天执行删除。先连续观察数个调度周期的预演日志,确认排序、目录层级、时区和年龄阈值都符合预期,再启用 --execute。
五、HDFS 已删除,为什么资源中心仍然显示
本次操作中,清理任务日志显示成功:
MODE=EXECUTE
SUMMARY mode=execute candidates=4
随后直接检查 HDFS,相关路径也已经不存在,但 DolphinScheduler 资源中心界面仍能看到旧目录。
这并不必然说明 HDFS 删除失败,而是因为两者不是同一个层面的“目录视图”:
- HDFS NameNode 管理文件和目录是否真实存在;
- DolphinScheduler 资源中心是管理面,用来维护脚本、SQL、JAR 等调度资源;
- 部分 DolphinScheduler 版本还会维护资源元数据,页面也可能存在缓存;
- 在 DolphinScheduler 之外直接修改资源中心对应的 HDFS 路径,可能绕过资源中心自身的删除和元数据更新流程。
较早版本的 DolphinScheduler 讨论中也说明,资源中心会维护资源记录;外部直接修改 HDFS、S3 或本地文件系统时,资源信息并不会实时同步。不同版本的具体实现可能变化,因此应以实际版本的表结构、接口和缓存机制为准。
HDFS 是否存在,以命令结果为准
可以在 DolphinScheduler 新建一个只读 Shell 任务:
#!/bin/bash
set -eu
ROOT="hdfs:///flink/checkpoints/test"
echo "=== list root ==="
hdfs dfs -ls "$ROOT"
echo "=== list all descendants ==="
hdfs dfs -ls -R "$ROOT"
检查某一个旧目录是否还存在,可以使用:
TARGET="hdfs:///flink/checkpoints/test/flinkjob1/old-job-id"
if hdfs dfs -test -e "$TARGET"; then
echo "EXISTS: $TARGET"
exit 1
else
echo "NOT_FOUND: $TARGET"
fi
如果输出 NOT_FOUND,说明原路径已经从 HDFS 命名空间移除。若 HDFS 启用了 Trash,数据还可能暂时位于回收站并继续占用空间,但原目录已经不存在。
不要把 Checkpoint 放在资源中心目录下
测试过程中曾使用:
hdfs://172.30.50.12:8020/dolphinscheduler/hdfs/resources/ckpt_storage
这会把 Flink 运行状态放进 DolphinScheduler 资源中心的命名空间。结果是:Flink 和清理脚本把它当普通 HDFS 目录,DolphinScheduler 却可能把它当受管理资源,最终产生显示不一致。
生产环境应拆开:
# DolphinScheduler 脚本、SQL、JAR
hdfs:///dolphinscheduler/resources/
# Flink 自动 Checkpoint
hdfs:///flink/checkpoints/<env>/<app>/
# Flink 人工 Savepoint
hdfs:///flink/savepoints/<env>/<app>/
对于已经放入资源中心的旧目录,应先迁移 Flink 配置,让新作业写入独立根目录。资源中心中的残留条目,应通过对应版本的 DolphinScheduler 页面或资源 API 清理,而不是直接删除其数据库记录。
六、生产级生命周期管理建议
“每个作业只保留最新目录”适合验证流程或空间告急时的受控清理,但它不是完整的生产策略。成熟环境通常会同时管理恢复目标、发布状态、数据所有权和存储成本。
1. 不只按修改时间判断
修改时间只能说明目录最近何时变化,不能证明它是否仍有恢复价值。更可靠的清理条件应全部满足:
- 旧 Job ID 已确认终止;
- 新版本已经稳定运行;
- 新作业已经成功生成自己的 Checkpoint;
- 回滚保护期已经结束;
- 该路径不在发布系统的保护清单中;
- 最近一次恢复演练确认恢复链路有效。
2. 保留“恢复代际”,而不只是保留一个目录
建议至少区分:
- 当前稳定版本的恢复点;
- 上一个可回滚版本的恢复点;
- 发布前创建的 Savepoint;
- 已过期、可以回收的历史状态。
具体保留时长应由业务恢复目标、状态大小、发布频率和 HDFS 容量决定。例如保留最近两个可回滚代际、历史状态保留 3~7 天,可以作为起点,但不是适用于所有公司的固定标准。
3. 建立恢复点登记表
发布平台或运维系统至少记录:
app
environment
flink_job_id
code_version
restore_path
restore_type
claim_mode
created_at
status
protected_until
清理器根据登记状态决策,比单纯扫描 HDFS 目录安全得多。
4. 正确处理 Savepoint Claim Mode
从 Savepoint 恢复时,状态文件的所有权会影响能否删除原路径:
NO_CLAIM:Flink 不接管原快照;通常应等新作业完成第一次成功的完整 Checkpoint 后,再删除原恢复点;CLAIM:Flink 接管该快照,可能继续复用其文件,不应由外部脚本随意删除;LEGACY:所有权语义不够清晰,新系统应避免依赖。
5. 使用两阶段删除
推荐流程是:
发现候选项 → 生成审计清单 → 进入隔离期或 HDFS Trash → 延迟物理回收
脚本默认预演,执行删除时也不使用 -skipTrash,就是为了保留恢复窗口。与此同时,需要监控 Trash 占用,并配置合理的 fs.trash.interval 和检查点周期;否则“删除成功”不等于磁盘空间立即释放。
6. 建立容量和失败监控
至少监控:
- 每个应用和环境的 Checkpoint 总容量;
- Checkpoint 成功率、时长和最近成功时间;
- 历史 Job ID 数量;
- 清理候选数、清理失败数和回收空间;
- HDFS 配额、Trash 占用和小文件数量;
- 最近一次恢复演练时间。
7. 最小权限和审计
清理账号只应拥有指定 Checkpoint 根目录的删除权限,不能拥有整个 HDFS 或 DolphinScheduler 资源中心的广泛权限。每次删除应记录任务实例、操作者、候选路径、保留原因、删除结果和恢复窗口。
七、结论
RETAIN_ON_CANCELLATION 没有失效,它只是把取消后的恢复能力和清理责任同时交给了用户。execution.checkpointing.num-retained 也不是历史 Job ID 的全局垃圾回收器。
这次问题最终揭示了三个必须分开的边界:
- Flink 运行状态:决定哪些恢复点仍被作业使用;
- HDFS 存储事实:决定目录是否真实存在、空间是否已经回收;
- DolphinScheduler 资源管理:维护调度资源及其页面视图,不应承担 Flink 状态存储职责。
测试阶段可以用“每个逻辑作业保留最新运行目录”的脚本验证流程;生产阶段则应把清理决策与发布记录、活跃 Job ID、Savepoint 所有权、回滚窗口和恢复演练结合起来。真正可靠的生命周期管理,不是删除得更快,而是能够证明:删掉的状态不会再被恢复链路需要。