七万号·数据平台实战手记

Back

课程讲解 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:

  1. 上游插入 Barrier,随数据流广播。
  2. 算子收到 Barrier 后,将当前状态快照写入状态后端。
  3. 所有算子完成快照后,本次 Checkpoint 完成。
Source ──[Barrier]──> Operator A ──[Barrier]──> Operator B
                          │                       │
                       快照                     快照
plaintext

Savepoint 与状态迁移#

  • Savepoint 是用户主动触发的完整快照,常用于升级、扩容。
  • 状态 Schema 变更需要 状态迁移器(State Migration) 配合。

故障恢复#

发生故障时,JobManager 会:

  1. 重启所有 Task。
  2. 从最近一次成功的 Checkpoint 恢复状态。
  3. 数据回放至 Barrier 位置,保证不丢不重。

小结#

状态与容错是流计算区别于批处理的根本。下一讲我们将动手实现一个端到端的实时数仓案例。

课程目录 · Flink 原理与实战

2/2 讲
  1. 01 Flink 运行时架构与执行图
  2. 02 Flink 状态管理与检查点容错
Flink 状态管理与检查点容错
https://realcpf.tech/journal/flink-internals-02
Author 刘佳成
Published at 2025年5月27日