Flink Checkpoint 与 Savepoint
概述
Flink 的精确一次语义靠 分布式快照(Checkpoint) 实现:周期性保存算子状态与数据流位置。本文讲清 Barrier 对齐机制、Checkpoint 执行流程、Exactly-Once 保证,以及 Savepoint 的区别与调优。
一、Checkpoint 是什么
周期性(如 60s)把全作业状态做一致性快照
状态(算子状态 + 键控状态)+ 输入位置(offset)
存到可靠存储(HDFS/S3)
故障时从最近快照恢复| 特性 | 说明 |
|---|---|
| 周期 | 可配置间隔 |
| 一致性 | 分布式快照(Barrier 对齐) |
| 存储 | 配置的 Checkpoint 目录 |
| 恢复 | 从最近成功快照 |
二、Barrier 对齐机制
2.1 什么是 Barrier
Barrier(屏障) 是 Checkpoint 的标记,随数据流注入:
数据流: [数据][数据][Barrier][数据][数据]
└ Checkpoint n 的边界2.2 对齐流程
1. Source 收到 Barrier,快照自己的 offset
2. Barrier 随数据一起流向下游算子
3. 算子收到所有输入流的 Barrier(对齐)
4. 算子做状态快照
5. 全部算子快照完成 → Checkpoint 成功多输入流对齐:
输入A:Barrier 已到 → 缓冲后续数据,等输入B
输入B:Barrier 到达 → 对齐完成,快照
对齐期间数据缓冲,保证一致性2.3 对齐的意义
| 对齐 | 保证 |
|---|---|
| 状态与数据位置一致 | 恢复后无重复/丢失 |
| Exactly-Once | 精确一次语义基础 |
三、Exactly-Once 语义
3.1 来源
数据源 offset 快照(如 Kafka offset)
+ 算子状态快照
+ 对齐的 Barrier
= 精确一次恢复3.2 端到端
| 端 | 机制 |
|---|---|
| Source | offset 快照,恢复时重放 |
| 计算 | 状态快照恢复 |
| Sink | 两阶段提交(事务) |
| Sink 语义 | 实现 |
|---|---|
| Exactly-Once | 两阶段提交 + 事务 |
| At-Least-Once | 直接写(可能重复) |
四、Checkpoint 配置
4.1 基础配置
java
// 开启 Checkpoint,每 60s
env.enableCheckpointing(60000);
// 语义
env.getCheckpointConfig()
.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// 超时
env.getCheckpointConfig().setCheckpointTimeout(60000);| 参数 | 说明 |
|---|---|
| 间隔 | Checkpoint 频率 |
| 模式 | EXACTLY_ONCE / AT_LEAST_ONCE |
| 超时 | 快照超时则失败 |
| 最小间隔 | 两次 Checkpoint 最小间隔 |
| 并发数 | 同时进行几个 Checkpoint |
4.2 配置项
| 配置 | 默认 | 说明 |
|---|---|---|
execution.checkpointing.interval | 0 | 间隔 |
execution.checkpointing.timeout | 10min | 超时 |
execution.checkpointing.min-pause | 0 | 最小间隔 |
execution.checkpointing.max-concurrent-checkpoints | 1 | 并发 |
state.checkpoints.num-retained | 1 | 保留数 |
state.checkpoints.dir | - | 存储目录 |
五、Savepoint
5.1 区别
| 对比 | Checkpoint | Savepoint |
|---|---|---|
| 触发 | 自动周期 | 手动 |
| 生命周期 | 自动清理 | 手动管理 |
| 用途 | 故障恢复 | 升级、迁移、调试 |
| 存储 | 配置目录 | 指定路径 |
| 兼容 | 版本内 | 跨版本/改作业 |
5.2 使用
bash
# 触发 Savepoint
flink savepoint <jobId> hdfs:///flink/savepoints
# 从 Savepoint 恢复
flink run -s hdfs:///flink/savepoints/savepoint-xxx job.jar
# 取消时保存
flink cancel -s hdfs:///flink/savepoints <jobId>| 场景 | 用 Savepoint |
|---|---|
| 升级 Flink 版本 | 是 |
| 修改作业逻辑 | 是(兼容性需验证) |
| 常规故障恢复 | 用 Checkpoint 即可 |
六、故障恢复流程
1. 检测到 Task 失败
2. 作业重启(按重启策略)
3. 从最近成功的 Checkpoint 加载状态
4. 数据源从快照 offset 重放
5. 恢复执行| 恢复信息 | 来源 |
|---|---|
| 状态 | Checkpoint 快照 |
| 输入位置 | Source offset 快照 |
| Sink 重放 | 两阶段提交/幂等 |
七、Checkpoint 调优
| 场景 | 优化 |
|---|---|
| 快照太大 | 增量 Checkpoint(RocksDB) |
| 快照太频繁 | 增大间隔、最小间隔 |
| 对齐延迟高 | 减少输入流或缓冲 |
| 恢复慢 | 状态小、增量快照 |
| 失败频繁 | 调大超时、检查背压 |
| 指标 | 说明 |
|---|---|
| 完成时间 | 快照耗时 |
| 失败次数 | 快照失败数 |
| 大小 | 快照数据量 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| Checkpoint 一直失败 | 超时/背压/存储问题,调整参数 |
| 恢复后数据重复 | Sink 未两阶段提交,或模式为至少一次 |
| 对齐期间卡顿 | 输入流不均衡,减少流或优化 |
| Savepoint 恢复失败 | 状态不兼容,检查并行度/算子变更 |
| 快照太大 | 增量 Checkpoint 或减少状态 |