Flink 源码与运行原理

流式计算 / 源码阅读 / LESSON 13

状态恢复与 Rescale 深度解析

解释状态恢复、Key Group、重分区与 Rescale 的实际代价,帮助设计可演进的作业。

阅读时间
60 分钟
学习路径
Flink 源码与运行原理
内容来源
Flink 深度笔记

预计阅读时间: 60 分钟 前置阅读: doc-03(Checkpoint), doc-04(StateBackend) 下一次阅读: doc-14(监控排障)


1. 恢复场景全景

场景触发方式恢复源并行度变化特殊处理
Job FailoverTask/Job 异常Checkpoint不变重新分配 Slot + 恢复 State
HA 切换JM 挂掉Checkpoint不变Dispatcher 恢复 JobGraph → 新 JobMaster
Savepoint Restoreflink run -s savepointSavepoint可变重新分配 KeyGroup + StateRedistribution
Rescale手动停止 + 改并行度 + 启动Savepoint变化KeyGroup → Subtask 映射改变
JM FailoverJM 进程崩溃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 重分配

状态恢复与 Rescale 深度解析 图 01

扩容 (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.javarestoreSavepoint()Savepoint 恢复入口
flink-runtime/.../state/KeyGroupRangeAssignment.javacomputeKeyGroupRangeForOperatorIndex()KeyGroup 分配计算
flink-runtime/.../state/StateAssignmentOperation.javaassignStates()State 重分配
flink-runtime/.../executiongraph/ExecutionVertex.javaresetForNewExecution(), restoreState()Execution 级恢复
flink-runtime/.../state/KeyGroupRange.javaKeyGroup 范围表示
flink-runtime/.../checkpoint/CheckpointStore.javagetLatestCheckpoint(), 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,确保安全回滚。


下一步