1. 为什么需要 Checkpoint?

Flink 的算子与函数天然可有状态,而状态是大多数复杂流式任务的地基。要让应用在故障后“像没出过事一样”继续跑,必须把状态和流位置定期快照下来——这就是 Checkpoint
有了它,Flink 能把作业恢复到失败前一致的处理位置与状态,提供等价于“无故障执行”的语义。

2. 启用前提(Prerequisites)

在生产环境中,务必准备两类“可回放/可持久化”的基建:

  • 可回放的数据源:Kafka / RabbitMQ / Kinesis / PubSub,或 HDFS/S3 等文件系统。
  • 可持久化的状态存储:分布式文件系统(HDFS、S3、GFS、NFS、Ceph…)。

小贴士:Checkpoint 目录应高可用所有 TM/JM 可读写、并具备足够吞吐

3. 快速上手:最小可用配置示例(Java)

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 1) 启用 checkpoint:每 1000 ms 触发一次
env.enableCheckpointing(1000);

// 2) Exactly-once(默认),更安全
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

// 3) 两次 checkpoint 之间至少停 500 ms,确保业务推进
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500);

// 4) 超时:60s
env.getCheckpointConfig().setCheckpointTimeout(60_000);

// 5) 连续失败容忍:2 次
env.getCheckpointConfig().setTolerableCheckpointFailureNumber(2);

// 6) 并发 checkpoint:1(与非对齐配套)
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);

// 7) 外部化 checkpoint:作业取消后保留
env.getCheckpointConfig().setExternalizedCheckpointRetention(
    ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);

// 8) 反压场景:启用非对齐 checkpoint(EXACTLY_ONCE + 并发=1才能生效)
env.getCheckpointConfig().enableUnalignedCheckpoints();

// 9) 指定持久化存储(推荐 HDFS/S3)
Configuration fsCfg = new Configuration();
fsCfg.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
fsCfg.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "hdfs:///flink/ckpt/my-job");
env.configure(fsCfg);

// 10) 允许图中部分任务完成后仍继续做 checkpoint(默认启用)
Configuration afterFinishCfg = new Configuration();
afterFinishCfg.set(CheckpointingOptions.ENABLE_CHECKPOINTS_AFTER_TASKS_FINISH, true);
env.configure(afterFinishCfg);

4. 核心概念与模式选择

4.1 处理语义

  • Exactly-once(推荐):端到端最强保障。对绝大多数业务正确性至关重要。
  • At-least-once:适合对延迟极敏感(毫秒级)的场景,允许偶发重复,需下游去重。

4.2 对齐 vs. 非对齐(Aligned vs. Unaligned)

  • 对齐:屏障跟随数据流,吞吐高时可能在反压下拉长时延。

  • 非对齐(Unaligned):把缓冲中的数据也纳入状态,允许屏障“超车”,显著缩短在反压下的 checkpoint 时长

    • 前提:EXACTLY_ONCEmaxConcurrentCheckpoints = 1
    • 可配超时降级:execution.checkpointing.aligned-checkpoint-timeout,超时后自动从对齐切换到非对齐。

4.3 并发与最小间隔

  • max-concurrent-checkpoints:并发进行的 checkpoint 个数。并发>1可提升恢复点密度,但会占用更多 IO/CPU。
  • min-pause:两次 checkpoint 的最小静默间隔。与并发=1组合,能确保有“业务推进窗口”。

4.4 外部化与保留策略

  • RETAIN_ON_CANCELLATION:取消/失败都保留,便于手动从最近成功点恢复
  • DELETE_ON_CANCELLATION:被取消则清理(更省空间)。

5. 图部分完成后的 Checkpoint 与算子生命周期

  1. 部分任务完成后继续做 Checkpoint(1.15 起默认启用):

    • 已完成的子任务不再参与后续 checkpoint。
    • 对自定义算子/UDF 至关重要:在 StreamOperator#finish()清空缓冲并写入必要的事务指针finish() 之后的 checkpoint 通常应是“空的”。
  2. UnionListState 特殊规则

    • 常用于存放“外部系统分区偏移的全局视图”(如 Kafka offsets)。
    • 若只完成了部分使用该状态的子任务,会导致全局视图丢失。Flink 要求要么全部完成、要么都未完成,避免不一致。
  3. 两阶段提交(2PC)算子

    • 所有算子到达 End-of-Data 后会立即触发最终 checkpoint并等待其成功,确保能提交所有记录。

6. 文件合并机制(Flink 1.20 实验特性)

多 checkpoint 的小文件洪泛会压垮底层文件系统的元数据服务。1.20 引入统一文件合并(MVP):

  • 开关:execution.checkpointing.file-merging.enabled = true

  • 合并粒度:

    • 共享作用域状态:子任务级合并
    • 私有作用域状态:TaskManager 级合并
  • 跨 checkpoint 合并:execution.checkpointing.file-merging.across-checkpoint-boundary = true

  • 文件池:

    • 非阻塞(默认):永远给你文件,但可能制造很多物理文件
    • 阻塞:小文件紧张时会等待回收
    • 配置:execution.checkpointing.file-merging.pool-blocking
  • 代价:会产生空间放大,用 execution.checkpointing.file-merging.max-space-amplification 做上限控制。

适用场景:大规模作业 / 高频 checkpoint / 后端是 S3/HDFS,显著缓解 NameNode/元数据压力

7. 生产就绪配置模板(flink-conf.yaml 片段)

# 基础语义
execution.checkpointing.mode: EXACTLY_ONCE
execution.checkpointing.interval: 10s           # 基础周期
execution.checkpointing.min-pause: 2s           # 两次之间留出业务推进窗口
execution.checkpointing.timeout: 2m
execution.checkpointing.tolerable-failed-checkpoints: 2
execution.checkpointing.max-concurrent-checkpoints: 1

# 存储位置(文件系统方案)
execution.checkpointing.storage: filesystem
execution.checkpointing.dir: hdfs:///flink/ckpt/my-job
execution.checkpointing.num-retained: 3

# 外部化策略
execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION

# 非对齐相关(反压友好)
execution.checkpointing.unaligned.enabled: true
execution.checkpointing.aligned-checkpoint-timeout: 5s  # 5s 内对齐不成则转非对齐

# 任务完成后仍继续 checkpoint
execution.checkpointing.checkpoints-after-tasks-finish: true

# 本地备份(大状态恢复更快,覆盖 keyed state)
execution.checkpointing.local-backup.enabled: true
execution.checkpointing.local-backup.dirs: /data/flink/localState

# 文件合并(1.20 实验)
execution.checkpointing.file-merging.enabled: true
execution.checkpointing.file-merging.max-file-size: 256mb
execution.checkpointing.file-merging.max-space-amplification: 2.0
execution.checkpointing.file-merging.across-checkpoint-boundary: true
execution.checkpointing.file-merging.pool-blocking: false

8. 端到端落地示例:Kafka → Flink → HDFS(Exactly-once)

要点串联:

  1. Kafka Source:开启可回放(保留期足够长,或 compact 策略)。
  2. Flink:Exactly-once + 非对齐 + 外部化 + HDFS/S3 checkpoint。
  3. Sink(HDFS):使用支持 2PC 的文件系统 sink(如 StreamingFileSink 或新版 FileSink),并保留最终 checkpoint。

示意代码:

env.enableCheckpointing(10_000); // 10s
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().enableUnalignedCheckpoints();
env.getCheckpointConfig().setExternalizedCheckpointRetention(
    ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);

Configuration fsCfg = new Configuration();
fsCfg.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
fsCfg.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "hdfs:///flink/ckpt/kafka2hdfs");
env.configure(fsCfg);

// Kafka source 配置省略(确保可回放、分区均衡)
// FileSink:滚动策略 + 2PC 提交(取决于 Flink 版本的 API)

9. 参数调优建议(经验法则)

  • 延迟敏感

    • AT_LEAST_ONCE + 高频 checkpoint(100~500ms)
    • max-concurrent-checkpoints 视资源决定,>1 可提升密度
    • 非对齐:如有反压,强烈建议打开
  • 吞吐优先 / 大状态

    • EXACTLY_ONCE + checkpoint 间隔拉长(10~60s)
    • 启用增量 checkpoint(后端支持时)
    • 开启本地备份文件合并(1.20+)
    • num-retained ≥ 2,便于回滚
  • 外部系统事务(MySQL、Kafka 事务 Sink):

    • 一定要等待最终 checkpoint完成(框架已自动触发),自定义算子确保在 finish() 正确处理事务指针

10. 常见坑与排查清单(Checklist)

  1. Checkpoint 超时

    • 现象:Checkpoint X expired before completing
    • 排查:存储写入慢?反压严重?提高 timeout、开启非对齐、检查网络与对象存储带宽/限流。
  2. 目录权限/可见性

    • 现象:Permission denied 或 JM/TM 路径不一致
    • 排查:统一 execution.checkpointing.dir,确保所有节点可读写、Kerberos/HDFS Token 正确。
  3. UnionListState 不一致

    • 现象:部分子任务完成后 checkpoint 一直不成功
    • 处理:确保使用该状态的子任务要么都完成要么都未完成;或调整任务并发与算子设计。
  4. 文件洪泛

    • 现象:NameNode/元数据压力高,小文件数量爆炸
    • 处理:1.20+ 启用文件合并;增大滚动策略阈值;调大 checkpoint 间隔。
  5. Exactly-once 无法达成

    • 现象:端到端仍有重复/丢数
    • 检查:Source/Sink 是否支持端到端语义(Kafka 事务 / FileSink 2PC);中间算子是否引入外部副作用。

11. 迭代作业(Iterative Jobs)的特别说明

  • Flink 对非迭代作业提供处理保证;对迭代作业默认不支持
  • 强制开启env.enableCheckpointing(interval, EXACTLY_ONCE, true)
  • 风险:循环边上的在途记录与相应状态变更,在故障时会丢失。谨慎评估业务可接受性。

12. FAQ

Q:我应该配“checkpoint 间隔”还是“最小间隔(min-pause)”?
A:在生产里更推荐“最小间隔”,它对后端偶发变慢更鲁棒,能稳定留出业务推进窗口。

Q:非对齐会增大状态体积吗?
A:会把缓冲数据也落到状态里,时延换空间。在高反压场景非常划算。

Q:保留多少个历史 checkpoint 合适?
A:建议 ≥2(通常 2~3 即可),既能回滚又不至于占用过多存储。

Q:外部化 checkpoint 选保留还是删除?
A:开发/测试环境建议删;生产建议保留,但要搭配生命周期管理/清理策略

13. 结语

Checkpoint 是 Flink 端到端一致性的关键“保险丝”。只要你在语义选择存储选型反压处理生命周期管理上把关到位,就能让作业在真实世界的抖动与故障面前“稳如老狗”。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐