课程讲解 Flink 原理与实战 · 第 2 讲
/ 共 2 讲
Flink 状态管理与检查点容错
状态是什么、存在哪里、怎么恢复?用一张图看懂 Checkpoint、Savepoint 与 Exactly-Once 的实现机制。
views
| comments
2 min
flink / 容错 上一讲我们理清了 Flink 的执行图。本讲进入流计算最核心的难题:状态。
状态分类#
- 算子状态(Operator State):绑定在算子实例上,典型应用是 Kafka Source 的位点。
- 键控状态(Keyed State):按 key 分区存储,支持 ValueState / ListState / MapState。
状态后端#
- RocksDBStateBackend:本地磁盘 + 内存缓存,适合大状态。
- HashMapStateBackend:纯内存,适合小状态、高吞吐场景。
Checkpoint 机制#
Flink 通过 Barrier 对齐实现 Exactly-Once:
- 上游插入 Barrier,随数据流广播。
- 算子收到 Barrier 后,将当前状态快照写入状态后端。
- 所有算子完成快照后,本次 Checkpoint 完成。
Source ──[Barrier]──> Operator A ──[Barrier]──> Operator B
│ │
快照 快照plaintextSavepoint 与状态迁移#
- Savepoint 是用户主动触发的完整快照,常用于升级、扩容。
- 状态 Schema 变更需要 状态迁移器(State Migration) 配合。
故障恢复#
发生故障时,JobManager 会:
- 重启所有 Task。
- 从最近一次成功的 Checkpoint 恢复状态。
- 数据回放至 Barrier 位置,保证不丢不重。
小结#
状态与容错是流计算区别于批处理的根本。下一讲我们将动手实现一个端到端的实时数仓案例。