Published on

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

Authors
  • avatar
    Name
    Charles Chen
    Twitter

在 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。

清理规则是:

  1. 遍历 Checkpoint 根目录下的每个逻辑作业;
  2. 读取每个作业下的直属子目录;
  3. 按 HDFS 修改时间排序;
  4. 始终保留最新的一个运行目录;
  5. 其余目录只有超过最小年龄后才成为候选项;
  6. 默认只打印,只有显式传入 --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 的全局垃圾回收器。

这次问题最终揭示了三个必须分开的边界:

  1. Flink 运行状态:决定哪些恢复点仍被作业使用;
  2. HDFS 存储事实:决定目录是否真实存在、空间是否已经回收;
  3. DolphinScheduler 资源管理:维护调度资源及其页面视图,不应承担 Flink 状态存储职责。

测试阶段可以用“每个逻辑作业保留最新运行目录”的脚本验证流程;生产阶段则应把清理决策与发布记录、活跃 Job ID、Savepoint 所有权、回滚窗口和恢复演练结合起来。真正可靠的生命周期管理,不是删除得更快,而是能够证明:删掉的状态不会再被恢复链路需要。

参考资料