- Published on
Flink 1.18 Checkpoint Essentials
- Authors

- Name
- Charles Chen
做 Flink 实时任务时,State、StateBackend、Checkpoint、CheckpointStorage 这几个概念非常容易混在一起。
尤其是从老版本 Flink 升级过来的开发者,经常会有这样的疑问:
- State 不就是 Checkpoint 吗?
- StateBackend 不就是决定 Checkpoint 存哪里吗?
EmbeddedRocksDBStateBackend(true)里的true又是什么意思?
本文以 Flink 1.18 为例,把这套机制从概念、选型到常用参数完整梳理一遍。
1. 先记住最重要的一句话
在 Flink 1.18 中:
State 是程序当前正在使用的数据;Checkpoint 是某一时刻 State 的持久化快照。
因此:
State ≠ Checkpoint
Checkpoint = State 在某个时间点的恢复快照
例如一个 Flink 算子维护:
ValueState<Long> count;
当前:
count = 128
此时 128 属于运行中的 State。
12:00 做了一次 Checkpoint:
Checkpoint 100
└── count = 128
程序继续运行,到 12:01:
count = 156
那么现在:
当前 State = 156
Checkpoint 100 中保存的 State = 128
如果程序此时发生故障,Flink 可以利用 Checkpoint 100 恢复:
Checkpoint 100
↓
恢复 State
↓
count = 128
↓
从对应的数据位置继续处理
这就是 Checkpoint 最核心的作用:故障恢复。
2. StateBackend 和 CheckpointStorage 到底是什么?
理解了 State 和 Checkpoint,另外两个概念就很简单了。可以记成:
StateBackend
↓
运行中的 State 怎么保存
CheckpointStorage
↓
Checkpoint 最终保存在哪里
整个关系大概是:
Flink Job
│
▼
Runtime State
│
StateBackend
│
┌─────────────┴─────────────┐
▼ ▼
HashMapStateBackend EmbeddedRocksDBStateBackend
JVM Heap RocksDB / Local Disk
│ │
└─────────────┬─────────────┘
│
Checkpoint
│
▼
CheckpointStorage
│
┌─────────────┴─────────────┐
▼ ▼
JobManager FileSystem
HDFS/S3/...
所以:
StateBackend 管“活着的 State”,CheckpointStorage 管“备份出来的 State”。
3. Flink 1.18 有哪些 StateBackend?
Flink 1.18 最主要的两个 StateBackend 是:
| StateBackend | State 主要在哪里 | 特点 |
|---|---|---|
HashMapStateBackend | JVM Heap | 快,但受内存限制 |
EmbeddedRocksDBStateBackend | RocksDB | 能支撑超大 State |
3.1 HashMapStateBackend
配置:
env.setStateBackend(
new HashMapStateBackend()
);
它把运行中的 State 放在 JVM Heap 中。例如:
TaskManager JVM
┌───────────────────────────┐
│ Heap │
│ │
│ user1001 → count=100 │
│ user1002 → count=200 │
│ user1003 → count=300 │
│ ... │
└───────────────────────────┘
最大的优势是访问 State 快,因为 State 就在 JVM Heap 中。
比较适合:
- 中小规模 State
- 低延迟任务
- 普通聚合
- 状态生命周期较短
- JVM 内存比较充裕
但它的问题也很明显:State 大小受到 JVM 内存限制。
如果一个 TaskManager 只有 8 GB Heap,却需要维护几十 GB State,显然不适合 HashMapStateBackend。
4. EmbeddedRocksDBStateBackend
配置:
env.setStateBackend(
new EmbeddedRocksDBStateBackend(true)
);
它使用嵌入式 RocksDB 保存 Keyed State。大致结构:
TaskManager
│
▼
RocksDB
┌──────────────┐
│ MemTable │
│ Block Cache │
│ SST Files │
└──────┬───────┘
│
▼
TaskManager 本地磁盘
最大的优势是 State 可以远大于 JVM Heap。例如:
TaskManager Heap
8 GB
State
80 GB
HashMapStateBackend 很难处理,而 RocksDB 可以利用本地 SSD / Disk 保存大量 State。
因此非常适合:
- 大规模 Keyed State
- 超大窗口
- 大量用户状态
- 长 TTL State
- 实时去重
- 大型 Join
- 实时用户画像
- TB 级 State
代价则是 RocksDB 涉及序列化、JNI、内存与磁盘访问等,相比纯 Heap State 通常有更多开销。
因此:
RocksDB 不是 HashMap 的“高级版”,而是用一定性能成本换取更强的 State 扩展能力。
5. EmbeddedRocksDBStateBackend(true) 中的 true 是什么?
来看:
new EmbeddedRocksDBStateBackend(true)
这里的 true 表示开启 Incremental Checkpoint,也就是增量 Checkpoint。
例如 RocksDB 当前有 100 GB State:
RocksDB
001.sst 20 GB
002.sst 20 GB
003.sst 20 GB
004.sst 20 GB
005.sst 20 GB
总计 ≈ 100 GB
第一次 Checkpoint 需要保存这些 State。之后 RocksDB 又产生:
006.sst
007.sst
如果开启 Incremental Checkpoint,下一次 Checkpoint 可以复用之前已经持久化的 SST 文件,只上传新产生、需要更新的文件。
因此它并不是简单地检查哪些 key/value 改了,而是充分利用 RocksDB 的 SST 文件机制进行增量快照。
对于几十 GB、几百 GB 乃至 TB 级 State,这个能力非常重要。
6. CheckpointStorage 有哪些?
Flink 1.18 主要关注两个:
CheckpointStorage
│
├── JobManagerCheckpointStorage
│
└── FileSystemCheckpointStorage
6.1 JobManagerCheckpointStorage
Checkpoint 数据主要由 JobManager 内存承担:
env.getCheckpointConfig()
.setCheckpointStorage(
new JobManagerCheckpointStorage()
);
适合:
- 本地开发
- 测试
- Demo
- 极小 State
大型生产任务一般不会把它作为首选。
6.2 FileSystemCheckpointStorage
生产环境更加常见:
env.getCheckpointConfig()
.setCheckpointStorage(
new FileSystemCheckpointStorage(
"hdfs:///flink/checkpoints/my-job"
)
);
Checkpoint 可以保存到:
- HDFS
- S3
- GCS
- 其他 Flink 支持的持久化文件系统
这也是生产环境最常见的 CheckpointStorage。
7. 生产环境 StateBackend 怎么选?
可以简单理解:
State 大吗?
│
┌───────┴───────┐
│ │
不大 很大
│ │
▼ ▼
HashMap RocksDB
│
JVM Heap 能否
稳定装下?
│
能 → HashMap
悬 → RocksDB
实际应该关注的不是整个 Job 的 State,而是:
总 State 大小
+ 并行度
+ 每个 TaskManager / Subtask 承担多少 State
+ State 是否持续增长
+ 可用内存
例如:
总 State = 100 GB
并行度 = 100
平均 ≈ 1 GB / subtask
未必必须使用 RocksDB。
但如果:
总 State = 100 GB
并行度 = 5
平均 ≈ 20 GB / subtask
就应该认真考虑 RocksDB。
8. 最常见的两种生产组合
8.1 中小 State
env.setStateBackend(
new HashMapStateBackend()
);
env.getCheckpointConfig().setCheckpointStorage(
new FileSystemCheckpointStorage(
"hdfs:///flink/checkpoints/my-job"
)
);
工作过程:
运行:
State
↓
JVM Heap
Checkpoint:
JVM State
↓
Snapshot
↓
HDFS
8.2 大 State
env.setStateBackend(
new EmbeddedRocksDBStateBackend(true)
);
env.getCheckpointConfig().setCheckpointStorage(
new FileSystemCheckpointStorage(
"hdfs:///flink/checkpoints/my-job"
)
);
工作过程:
运行:
State
↓
RocksDB
↓
TaskManager Local SSD
Checkpoint:
RocksDB Snapshot
↓
HDFS
这也是理解 StateBackend 和 CheckpointStorage 最直观的方式。
9. Checkpoint 最重要的参数
理解完架构,再来看真正影响生产稳定性的参数。
9.1 Checkpoint Interval
开启 Checkpoint:
env.enableCheckpointing(
60_000,
CheckpointingMode.EXACTLY_ONCE
);
这里 60_000 ms = 60 秒,表示 Checkpoint 的触发间隔。它控制的是:希望多久触发一次 Checkpoint。
注意,它不意味着每 60 秒必然产生一个成功的 Checkpoint。如果前一个 Checkpoint 很慢,还会受到并发数、最小暂停时间等参数限制。
10. CheckpointingMode
一般配置:
env.getCheckpointConfig()
.setCheckpointingMode(
CheckpointingMode.EXACTLY_ONCE
);
主要有:
EXACTLY_ONCEAT_LEAST_ONCE
绝大多数需要一致性的生产 Stateful Streaming Job,优先使用 EXACTLY_ONCE。
需要注意,Checkpoint 的 Exactly Once 与端到端 Exactly Once 不是完全等价的概念。
如果是:
Kafka → Flink → MySQL
最终能否做到端到端 Exactly Once,还取决于 Source、Sink 以及外部系统是否配合相应的一致性机制。
11. Checkpoint Timeout
config.setCheckpointTimeout(
10 * 60 * 1000
);
表示一次 Checkpoint 最长允许执行多久。
例如:
12:00
CKPT 100 START
│
│
│ 10 min
▼
12:10
仍未完成
│
▼
TIMEOUT
所以:
Checkpoint Interval
= 多久尝试做一次
Checkpoint Timeout
= 一次最多允许做多久
两者完全不同。
12. MinPauseBetweenCheckpoints
配置:
config.setMinPauseBetweenCheckpoints(
30_000
);
表示上一个 Checkpoint 完成之后,至少休息多久才能开始下一个。
例如:
CKPT 100
│
│ 执行 40s
▼
完成
│
│ 休息至少 30s
▼
CKPT 101
为什么需要这个参数?因为 Checkpoint 会产生:
- CPU
- 序列化
- 本地磁盘 IO
- 网络 IO
- HDFS / S3 IO
如果任务:
CKPT
↓
刚完成
↓
马上 CKPT
↓
刚完成
↓
又 CKPT
整个系统可能长期处于 Checkpoint 压力中。因此这个参数对于大状态任务尤其有价值。
13. MaxConcurrentCheckpoints
config.setMaxConcurrentCheckpoints(1);
表示最多允许几个 Checkpoint 同时执行。
设置为 1 意味着:
CKPT 100
███████████████
CKPT 101
█████████████
而不会大量重叠:
CKPT 100
████████████████
CKPT 101
███████████████
生产环境设置为 1 是一个非常常见的起点。
14. Interval、MinPause 和 Concurrent 要一起看
这是三个最容易混淆的参数。
假设:
env.enableCheckpointing(60_000);
config.setMinPauseBetweenCheckpoints(30_000);
config.setMaxConcurrentCheckpoints(1);
分别控制:
| 参数 | 含义 |
|---|---|
| Interval | 希望多久触发一次 |
| MinPause | 上次完成后至少休息多久 |
| MaxConcurrent | 最多允许几个同时执行 |
因此 Checkpoint 调度并不是简单的“每 60 秒无脑执行一次”,而是这些条件共同作用的结果。
15. TolerableCheckpointFailureNumber
配置:
config.setTolerableCheckpointFailureNumber(3);
表示可以容忍一定数量的 Checkpoint Failure。例如:
CKPT 100 SUCCESS
CKPT 101 FAILED
CKPT 102 FAILED
CKPT 103 FAILED
...
这个参数的意义在于:HDFS、S3、网络等偶尔出现瞬时异常时,不一定需要立即把整个 Job 干掉。
但也不能无限容忍。否则可能出现非常危险的情况:
Job Status:
RUNNING ✓
最近一次成功 Checkpoint:
8 小时前 ✗
看起来 Job 正常运行,实际上容灾能力已经严重下降。
16. Externalized Checkpoint
生产环境另一个非常重要的配置:
config.setExternalizedCheckpointCleanup(
ExternalizedCheckpointCleanup
.RETAIN_ON_CANCELLATION
);
意思是:Job 被 Cancel 后,保留 Externalized Checkpoint。
例如:
Job Running
CKPT 100
CKPT 101
CKPT 102
│
▼
Cancel
│
▼
Job 停止
但是 CKPT 仍然保留
以后可以利用它恢复。这种配置对于生产任务非常实用,但需要注意:保留下来的数据也需要进行生命周期和存储空间管理。
17. FileStateSizeThreshold
FileSystemCheckpointStorage 还有一个容易忽略的参数:
new FileSystemCheckpointStorage(
checkpointPath,
fileStateSizeThreshold
);
它控制:小块 State 是 inline 到 Checkpoint metadata,还是单独生成文件。
如果存在大量非常小的 State:
20 Byte
50 Byte
100 Byte
...
全部生成独立文件可能导致:
checkpoint/
state-000001
state-000002
state-000003
state-000004
state-000005
...
产生大量小文件。Threshold 可以用于控制小 State 的存储方式。
一般来说,没有明确的小文件问题时,不建议为了“优化”而随意调整这个参数。
18. Unaligned Checkpoint
这是处理反压问题时非常重要的机制:
config.enableUnalignedCheckpoints();
正常 Aligned Checkpoint 需要处理 Checkpoint Barrier。如果任务出现严重 Backpressure:
Source
│
▼
Operator A
│
│ ████████████
│ 严重 Backpressure
▼
Operator B
Barrier 可能长时间无法完成对齐。
于是:
正常:
Checkpoint = 10 秒
严重反压:
Checkpoint = 1 分钟
3 分钟
5 分钟
Timeout
Unaligned Checkpoint 的重要思想是:不再等待所有 in-flight 数据完成传统意义上的对齐,而是把相关网络中的 in-flight data 也纳入快照。
优点:严重反压时,Checkpoint 时间更加可控。
代价之一:Checkpoint Size 可能增大。
因此:
Unaligned Checkpoint 主要解决的是反压导致的 Barrier Alignment 问题,而不是“让所有 Checkpoint 都变快”的万能开关。
19. AlignedCheckpointTimeout
还可以采用一种折中方案:
config.setAlignedCheckpointTimeout(...);
思路:
开始 CKPT
│
▼
先尝试 Aligned
│
│
│ 等待一定时间
▼
Alignment 还是很慢
│
▼
切换 Unaligned
这特别适合平时没有严重反压、偶尔出现反压尖峰的任务。
20. 为什么不建议生产使用 file:///tmp/flink/
开发环境经常看到:
new FileSystemCheckpointStorage(
"file:///tmp/flink/"
);
本地测试没有问题,但生产集群需要谨慎。
因为 file:// 通常代表节点本地文件系统。例如:
TaskManager 1
/local/tmp/flink
TaskManager 2
/local/tmp/flink
这两个目录不是同一个共享目录。如果机器、容器或者 Pod 消失,本地文件也可能一起消失。
Checkpoint 最重要的目的恰恰是:当前执行节点出问题之后还能恢复。
所以生产环境通常应该考虑:
- HDFS
- S3
- GCS
- 其他可靠、持久、集群可访问的存储
21. 一套 Flink 1.18 Checkpoint 配置模板
把前面的内容组合起来,一个 RocksDB Stateful Streaming Job 可以从下面这样的配置开始:
// ==========================================
// 1. 每 60 秒尝试进行一次 Checkpoint
// ==========================================
env.enableCheckpointing(
60_000,
CheckpointingMode.EXACTLY_ONCE
);
CheckpointConfig ckpt =
env.getCheckpointConfig();
// ==========================================
// 2. 单次 Checkpoint 最长 10 分钟
// ==========================================
ckpt.setCheckpointTimeout(
10 * 60 * 1000
);
// ==========================================
// 3. 两个 Checkpoint 至少暂停 30 秒
// ==========================================
ckpt.setMinPauseBetweenCheckpoints(
30_000
);
// ==========================================
// 4. 最多同时执行一个 Checkpoint
// ==========================================
ckpt.setMaxConcurrentCheckpoints(1);
// ==========================================
// 5. Checkpoint Failure 容忍策略
// ==========================================
ckpt.setTolerableCheckpointFailureNumber(3);
// ==========================================
// 6. Cancel Job 后保留 Externalized CKPT
// ==========================================
ckpt.setExternalizedCheckpointCleanup(
ExternalizedCheckpointCleanup
.RETAIN_ON_CANCELLATION
);
// ==========================================
// 7. RocksDB StateBackend
// true = Incremental Checkpoint
// ==========================================
env.setStateBackend(
new EmbeddedRocksDBStateBackend(true)
);
// ==========================================
// 8. Checkpoint 持久化到 HDFS
// ==========================================
ckpt.setCheckpointStorage(
new FileSystemCheckpointStorage(
"hdfs:///flink/checkpoints/my-job"
)
);
注意:这是一套便于理解的起始模板,而不是所有生产任务都应该照抄的“最佳参数”。
真正的参数应该根据 State 大小、Checkpoint Duration、反压、磁盘 IO、网络以及 HDFS/S3 性能进行调整。
22. 最值得记住的 Checkpoint 参数
最后浓缩成一张表:
| 参数 | 作用 | 常见思路 |
|---|---|---|
enableCheckpointing(interval) | 多久尝试触发一次 CKPT | 根据恢复目标和开销设置 |
CheckpointingMode | 一致性语义 | 通常 EXACTLY_ONCE |
checkpointTimeout | CKPT 最长允许多久 | 应明显高于正常 CKPT 时间 |
minPauseBetweenCheckpoints | 两次 CKPT 之间最小休息时间 | 防止持续 CKPT |
maxConcurrentCheckpoints | 最多几个 CKPT 并发 | 常从 1 开始 |
tolerableCheckpointFailureNumber | 容忍多少 CKPT Failure | 根据 SLA 设置 |
RETAIN_ON_CANCELLATION | Cancel 后是否保留 | 生产常考虑 |
RocksDB(true) | Incremental CKPT | 大 State 常用 |
CheckpointStorage | CKPT 保存位置 | 生产用可靠持久存储 |
fileStateSizeThreshold | 小 State 是否 inline | 通常保持默认 |
UnalignedCheckpoint | 缓解反压导致的 Alignment 问题 | 有对应问题再使用 |
AlignedCheckpointTimeout | Aligned 多久后转 Unaligned | 适合间歇性反压 |
23. 出现 Checkpoint 慢,应该看什么?
这是生产排障中特别重要的一点:Checkpoint 慢,不要第一反应就是把 checkpointTimeout 从 10 分钟改成 30 分钟。
Timeout 只是允许它慢得更久,并没有解决为什么慢。
应该进入 Flink Web UI 观察 Checkpoint 指标,重点关注:
- Checkpoint Duration
- Checkpoint Size
- Start Delay
- Alignment Duration
- Sync Duration
- Async Duration
然后大致按照这个思路排查:
Checkpoint 慢
│
├── Alignment 很慢?
│ │
│ └── 看 Backpressure
│ 看 Unaligned CKPT
│
├── Async Snapshot 很慢?
│ │
│ ├── RocksDB?
│ ├── 本地磁盘?
│ └── HDFS/S3?
│
├── Checkpoint Size 特别大?
│ │
│ ├── State 是否暴涨?
│ ├── TTL 是否合理?
│ └── Incremental CKPT?
│
└── Start Delay 很大?
│
└── Barrier 传播 / 上游反压
这比简单增加 Timeout 更有意义。
总结
Flink Checkpoint 体系看起来参数很多,但核心其实可以浓缩成四个概念:
State
│
│ 当前运行的数据
▼
StateBackend
│
│ 决定 State 怎么存
│
├── HashMap
└── RocksDB
│
│ Snapshot
▼
Checkpoint
│
│ 持久化
▼
CheckpointStorage
│
├── JobManager
└── FileSystem
│
├── HDFS
├── S3
└── ...
选型上则可以先记住:
中小 State + 追求访问性能
↓
HashMapStateBackend
+
FileSystemCheckpointStorage
大 State / State 可能超过内存
↓
EmbeddedRocksDBStateBackend(true)
+
FileSystemCheckpointStorage
最后再记住六个最核心的调优问题:
Interval
→ 多久做一次?
Timeout
→ 一次最多做多久?
MinPause
→ 做完至少休息多久?
MaxConcurrent
→ 最多同时做几个?
Incremental
→ RocksDB 是否增量保存?
Unaligned
→ 严重反压时 Barrier 怎么处理?
掌握这套关系之后,Flink Checkpoint 的配置就不再是“背参数”,而是在控制三件事情:多久备份一次、备份如何执行、备份最终放在哪里。