预计阅读时间: 60 分钟 前置阅读: doc-03(Checkpoint), doc-04(StateBackend) 下一次阅读: doc-14(监控排障)
1. 恢复场景全景
| 场景 | 触发方式 | 恢复源 | 并行度变化 | 特殊处理 |
|---|---|---|---|---|
| Job Failover | Task/Job 异常 | Checkpoint | 不变 | 重新分配 Slot + 恢复 State |
| HA 切换 | JM 挂掉 | Checkpoint | 不变 | Dispatcher 恢复 JobGraph → 新 JobMaster |
| Savepoint Restore | flink run -s savepoint | Savepoint | 可变 | 重新分配 KeyGroup + StateRedistribution |
| Rescale | 手动停止 + 改并行度 + 启动 | Savepoint | 变化 | KeyGroup → Subtask 映射改变 |
| JM Failover | JM 进程崩溃 | HA Metadata | 不变 | 新的 JM 接管 ExecutionGraph |
2. KeyGroup — Rescale 的核心
2.1 为什么需要 KeyGroup?
直接映射: Key → hash(key) % parallelism → Subtask
问题: 改变 parallelism 后所有 Key 的映射都变了 → State 全部失效!
KeyGroup 映射: Key → hash(key) % maxParallelism → KeyGroup → Subtask
优势: KeyGroup 映射在 Key→KeyGroup 层不变,只重新分配 KeyGroup→Subtask
2.2 KeyGroup 分配算法
// KeyGroupRangeAssignment.java — 核心算法
public static int computeKeyGroupForKeyHash(int keyHash, int maxParallelism) {
return MathUtils.murmurHash(keyHash) % maxParallelism;
}
public static KeyGroupRange computeKeyGroupRangeForOperatorIndex(
int maxParallelism, int parallelism, int operatorIndex) {
int start = operatorIndex * maxParallelism / parallelism;
int end = (operatorIndex + 1) * maxParallelism / parallelism - 1;
return new KeyGroupRange(start, end);
}
// 例: maxParallelism=128, parallelism=4
// Subtask 0: KeyGroup [0, 31] ← ceil(128/4) = 32 个
// Subtask 1: KeyGroup [32, 63]
// Subtask 2: KeyGroup [64, 95]
// Subtask 3: KeyGroup [96, 127]
//
// Rescale → parallelism=8:
// Subtask 0: KeyGroup [0, 15] ← ceil(128/8) = 16 个
// Subtask 1: KeyGroup [16, 31]
// ...
// Subtask 7: KeyGroup [112, 127]
2.3 MaxParallelism 的选择
MaxParallelism 推荐: 2^n, 且 128 ≤ maxParallelism ≤ 32768
示例:
预期最大并行度=16 → MaxParallelism=128 (安全余量 8x)
预期最大并行度=100 → MaxParallelism=256 (安全余量 2.5x)
预期最大并行度=500 → MaxParallelism=1024 (安全余量 2x)
不可以改 MaxParallelism → 修改意味着 Key→KeyGroup 映射全部改变
→ 存量 State 全部失效!
3. Savepoint Restore 流程
3.1 整体流程
// ExecutionGraph.java — Savepoint 恢复
public void restoreSavepoint(SavepointRestoreSettings settings) {
// 1. 加载 Savepoint
CompletedCheckpoint savepoint =
checkpointStore.getCheckpointByExternalPointer(settings.getRestorePath());
// 2. 解析 Operator State
Map<OperatorID, OperatorState> operatorStates = savepoint.getOperatorStates();
// OperatorState = Map<KeyGroup, KeyedStateHandle>
// 3. 如果并行度变化 → 重新计算 KeyGroup→Subtask 映射
// 4. 状态重分配
StateAssignmentOperation.assignStates(
operatorStates, newExecutionVertices, maxParallelism);
}
// StateAssignmentOperation.java — 状态重分配核心
public static Map<ExecutionAttemptID, List<StateHandle>> assignStates(
Map<OperatorID, OperatorState> oldStates,
Collection<ExecutionVertex> newVertices,
int maxParallelism) {
Map<ExecutionAttemptID, List<StateHandle>> assigned = new HashMap<>();
for (ExecutionVertex vertex : newVertices) {
// 计算这个 Subtask 负责的 KeyGroup 范围
KeyGroupRange range = KeyGroupRangeAssignment
.computeKeyGroupRangeForOperatorIndex(
maxParallelism, vertex.getParallelism(), vertex.getSubtaskIndex());
// 从旧 OperatorState 中提取对应 KeyGroup 的 StateHandle
OperatorState oldState = oldStates.get(vertex.getOperatorID());
List<StateHandle> handles = oldState.getStateHandlesByKeyGroupRange(range);
// 多个 old Subtask 的 StateHandle 可能分配到同 1 个 new Subtask
// 1 个 old Subtask 的 StateHandle 可能拆分到多个 new Subtask
assigned.put(vertex.getExecutionId(), handles);
}
return assigned;
}
3.2 并行度变化的 State 重分配

扩容 (scale-up):每个 old Subtask 的 State 拆分到多个 new Subtask
- old Subtask 0 (KeyGroup [0,31]) → new Subtask 0 ([0,15]) + new Subtask 1 ([16,31])
- StateHandle 拆开可能涉及网络传输(跨 TM)
缩容 (scale-down):多个 old Subtask 的 State 合并到 1 个 new Subtask
- old Subtask 0 ([0,31]) + old Subtask 1 ([32,63]) → new Subtask 0 ([0,63])
- 多个 StateHandle 合并可能需要下载 + 合并(RocksDB 的 SST 文件)
4. 恢复性能优化
4.1 增量 Checkpoint 的恢复
增量 Checkpoint 恢复流程:
1. 从 CompletedCheckpoint 获取 StateHandle
2. IncrementalRemoteKeyedStateHandle 包含:
- sharedState: 共享的 SST 文件(来自之前 Checkpoint)
- privateState: 本次独有的 SST 文件
3. 恢复时: 下载 ALL SST 文件(共享 + 私有)
4. 多文件并行下载 → 恢复速度受限于网络带宽 / 文件数
4.2 Local Recovery
Local Recovery:
Checkpoint → 同时写 HDFS/S3 + 本地磁盘
TaskManager 重启(进程重启但磁盘完好)
→ 优先从本地恢复 → 跳过网络下载
TaskManager 挂掉(节点丢失)
→ 从远程 HDFS/S3 恢复
4.3 Region Recovery
Region Recovery (FLIP-311):
传统的全局恢复: 1 个 Task 失败 → 整个 Job 重启
Region Recovery: 1 个 Task 失败 → 只重启受影响的 Region
Region = PipelinedRegion (FORWARD 连接的一组算子)
跨 Region 的关系: HASH/RESCALE/BROADCAST (Blocking Shuffle)
Region Recovery 的条件:
- 使用 Blocking Shuffle 或 Pipelined Bounded Shuffle
- 失数 Region 可以独立重启 + 从上一个 Region 的持久化数据恢复
5. 源码导航
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-runtime/.../checkpoint/CheckpointCoordinator.java | restoreSavepoint() | Savepoint 恢复入口 |
flink-runtime/.../state/KeyGroupRangeAssignment.java | computeKeyGroupRangeForOperatorIndex() | KeyGroup 分配计算 |
flink-runtime/.../state/StateAssignmentOperation.java | assignStates() | State 重分配 |
flink-runtime/.../executiongraph/ExecutionVertex.java | resetForNewExecution(), restoreState() | Execution 级恢复 |
flink-runtime/.../state/KeyGroupRange.java | — | KeyGroup 范围表示 |
flink-runtime/.../checkpoint/CheckpointStore.java | getLatestCheckpoint(), getCheckpointByExternalPointer() | Checkpoint/Savepoint 检索 |
6. 常见问题 / 面试题
Q1: MaxParallelism 为什么创建后不能改?
A: MaxParallelism 决定了 KeyGroup 总量,而 Key→KeyGroup 映射(hash(key) % maxParallelism)绑定 KeyGroup 总量。如果改了 maxParallelism,Key→KeyGroup 映射全部改变,存量 State(以 KeyGroup 为单位存储)无法在新映射下恢复。
Q2: Rescale 时会丢失数据吗?
A: 不丢失。状态按 KeyGroup 重分配到新 Subtask,所有 State 在新 SubTask 恢复——只是可能涉及网络传输和下载。但 Rescale 过程中有一个不可用窗口(恢复阶段),这段时间不处理新数据。
Q3: Local Recovery 和 Remote Recovery 各自的适用场景?
A: Local Recovery (本地磁盘) 适合 TM 进程重启(OOM 后被 YARN/K8s 重启)→ 磁盘完好 → 秒级恢复。Remote Recovery (HDFS/S3) 适合 TM 节点丢失(物理机器挂掉)→ 从远程恢复 → 受网络带宽限制。推荐两者都启用。
Q4: Savepoint 要如何生成才能跨版本恢复?
A: (1) 使用 Canonical 格式(--type canonical),保证自包含;(2) 不要启用增量 Savepoint;(3) 保留用户代码的兼容性(State 序列化器不能 breaking change);(4) 建议在升级 Flink 版本前生成 Savepoint,确保安全回滚。