Flink 源码与运行原理

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

Job 提交与调度链路——从 JAR/SQL 到分布式 Task 的全源码穿越

追踪一个 JAR 或 SQL 作业如何穿过 Dispatcher、JobMaster 和 Scheduler,最终落成分布式 Task。

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

预计阅读时间: 120 分钟 前置阅读: doc-00——架构概览与代码地图 下一次阅读: doc-02(流数据链路), doc-11(ResourceManager)


Job 提交与调度链路——从 JAR/SQL 到分布式 Task 的全源码穿越 图 01

1. Client 层—JAR/SQL 到 JobGraph

1.1 调用链

CLI: flink run job.jar
  → CliFrontend.main()
    → CliFrontend.run()
      → PackagedProgram.invokeInteractiveModeForExecution()
        → StreamGraphGenerator.generate()          // DataStream API → StreamGraph
          → StreamingJobGraphGenerator.createJobGraph()  // StreamGraph → JobGraph
            → [Operator Chain 优化] → [SlotSharingGroup 分组]
              → JobGraph (JobVertex[] + IntermediateDataSet[])

SQL: INSERT INTO sink SELECT ... FROM source
  → TableEnvironmentImpl.execute()
    → Planner.optimize()                    // Table/SQL → Transformation
      → translateToJobGraph()               // Transformation → JobGraph

1.2 StreamGraph → JobGraph 的转换细节

StreamGraph
├── StreamNode 0 (Source: Kafka Consumer)    ← type: Source
├── StreamNode 1 (Map: JSON Parse)            ← type: OneInput
├── StreamNode 2 (KeyBy: user_id)             ← type: KeyBy ← 引入 Shuffle!
├── StreamNode 3 (Window: 1min)              ← type: Window
├── StreamNode 4 (Agg: COUNT)                ← type: OneInput
└── StreamNode 5 (Sink: JDBC)                ← type: Sink

    [Operator Chain 优化后]

JobGraph
├── JobVertex 0: Source→Map (chained)        ← 并行度=8, SlotSharingGroup=default
├── JobVertex 1: Window→Agg (chained)        ← 并行度=4, SlotSharingGroup=default
└── JobVertex 2: Sink                        ← 并行度=4, SlotSharingGroup=default
    IntermediateDataSet: Source→Window > (HASH by user_id)  ← 跨网络 Shuffle
    IntermediateDataSet: Window→Sink > (FORWARD)             ← 同 Task 内

Operator Chain 的判定条件(源码级)

// StreamingJobGraphGenerator.java — 链化条件检查
private static boolean isChainable(StreamNode upstream, StreamNode downstream) {
    // 1. 并行度相同
    if (upstream.getParallelism() != downstream.getParallelism()) return false;

    // 2. 相同的 SlotSharingGroup
    if (!upstream.getSlotSharingGroup().equals(downstream.getSlotSharingGroup()))
        return false;

    // 3. 数据分发为 FORWARD(不能有 SHUFFLE/REBALANCE/RESCALE)
    // KeyBy 产生了 HASH 分发 → 触发链断开
    if (upstream.getOutEdges().stream()
            .anyMatch(edge -> edge.getPartitioner() != ForwardPartitioner.INSTANCE))
        return false;

    // 4. 无 Chain 禁用标记
    if (upstream.getChainingStrategy() == ChainingStrategy.NEVER) return false;

    return true;
}

1.3 StreamGraph vs JobGraph

维度StreamGraphJobGraph
粒度StreamNode (每个转换算子一个)JobVertex (多个算子合并后一个)
StreamEdge (算子间数据流)IntermediateDataSet → JobEdge (数据依赖)
优化无优化Operator Chain + SlotSharingGroup 优化
生成者Client (用户代码)StreamJobGraphGenerator (优化器)
内容用户级的逻辑表达为调度准备的执行图表示

2. Dispatcher—接收 JobGraph, 启动 JobMaster

2.1 Dispatcher 源码结构

// Dispatcher.java — Job 提交入口
public class Dispatcher extends FencedRpcEndpoint<DispatcherId> {
    // 已注册的 JobManagerRunner (运行中的 Job)
    private final JobManagerRunnerRegistry jobManagerRunnerRegistry;

    // 正在等待之前 Job 终止的 JobID 集合 (去重)
    private final Set<JobID> submittedAndWaitingTerminationJobIDs;

    // JobGraph 持久化存储 (HA)
    private final JobGraphStore jobGraphStore;

    // JobResult 存储 (全局去重)
    private final JobResultStore jobResultStore;
}

2.2 submitJob — 提交流程

// Dispatcher.java — submitJob 核心流程
public CompletableFuture<JobID> submitJob(JobGraph jobGraph, Time timeout) {
    // 1. 全局去重: 检查 Job 是否已是终端状态 (已完成/已取消/已失败)
    if (jobResultStore.hasJobResultEntryAsync(jobGraph.getJobID())) {
        // 如果是重复提交的已完成 Job → 直接返回 JobID (幂等)
        return CompletableFuture.completedFuture(jobGraph.getJobID());
    }

    // 2. 本地去重: 检查是否正在等待之前的 Job 终止
    if (submittedAndWaitingTerminationJobIDs.contains(jobGraph.getJobID())) {
        throw new DuplicateJobSubmissionException("Job already submitted");
    }

    // 3. 应用并行度覆盖(来自 flink-conf.yaml)
    applyParallelismOverrides(jobGraph);

    // 4. 提交
    return internalSubmitJob(jobGraph);
}

private CompletableFuture<JobID> internalSubmitJob(JobGraph jobGraph) {
    submittedAndWaitingTerminationJobIDs.add(jobGraph.getJobID());

    // 等待任何之前运行的 same JobID 的 JobMaster 优雅终止
    return waitForTerminatingJob(jobGraph.getJobID())
        .thenCompose(ignored -> persistAndRunJob(jobGraph));
}

private CompletableFuture<JobID> persistAndRunJob(JobGraph jobGraph) {
    // 1. 持久化 JobGraph 到 HA (ZooKeeper/HDFS)
    jobGraphStore.putJobGraph(jobGraph);

    // 2. 创建 JobManagerRunner (容器: 包含 JobMaster)
    JobManagerRunner runner = createJobMasterRunner(jobGraph);

    // 3. 启动 Runner (→ JobMaster.onStart() → startScheduling())
    return runJob(runner, ExecutionType.SUBMISSION);
}

2.3 恢复流程 (HA)

// Dispatcher.onStart() — 恢复之前持久化的 Job
public void onStart() {
    // ... 启动服务 ...

    // 恢复所有之前未完成的 Job
    Collection<JobGraph> recoveredJobs = jobGraphStore.getJobGraphs();
    for (JobGraph jobGraph : recoveredJobs) {
        // 创建 JobMasterRunner 并以 RECOVERY 模式启动
        runJob(createJobMasterRunner(jobGraph), ExecutionType.RECOVERY);
    }
}

3. JobMaster—JobGraph → ExecutionGraph → 调度

3.1 JobMaster 内部结构

// JobMaster.java — Per-Job Master
public class JobMaster extends FencedRpcEndpoint<JobMasterId> {
    // 调度器 (DefaultScheduler / AdaptiveScheduler)
    private final SchedulerNG schedulerNG;

    // SlotPool: JobMaster 端的 Slot 管理
    private final SlotPoolService slotPoolService;

    // 已注册的 TaskManager
    private final Map<ResourceID, TaskManagerRegistration> registeredTaskManagers;

    // 到 ResourceManager 的重试连接
    private ResourceManagerConnection resourceManagerConnection;

    // ShuffleMaster 连接
    private final PartitionTracker partitionTracker;
}

3.2 onStart → startScheduling

// JobMaster.java — 启动流程
public void onStart() {
    startJobExecution();
}

private void startJobExecution() {
    // 1. 向 ShuffleMaster 注册 Job
    shuffleMaster.registerJob(jobGraph.getJobID());

    // 2. 启动 JobMaster 服务
    startJobMasterServices();
    // a. 启动心跳管理器 (TM + RM)
    // b. 启动 SlotPoolService
    // c. 开始检索 RM Leader (等 RM 上线后连接)

    // 3. 开始调度!
    startScheduling();
}

private void startScheduling() {
    schedulerNG.startScheduling();  // → DefaultScheduler.startSchedulingInternal()
}

3.3 JobGraph → ExecutionGraph 的构建

JobGraph (Client 端构建)
  JobVertex[0]: "Source→Map"    parallelism=8
  JobVertex[1]: "Window→Agg"   parallelism=4
  JobVertex[2]: "Sink"           parallelism=4

          ↓ ExecutionGraphBuilder.buildGraph()

ExecutionGraph (JobMaster 端构建)
  ExecutionJobVertex[0] (Source→Map)
  ├── ExecutionVertex[0][0]  →  Execution(attempt=0, state=CREATED)
  ├── ExecutionVertex[0][1]  →  Execution(attempt=0, state=CREATED)
  ├── ... (8 个)
  │
  ExecutionJobVertex[1] (Window→Agg)
  ├── ExecutionVertex[1][0]  →  Execution(attempt=0, state=CREATED)
  ├── ... (4 个)
  │
  ExecutionJobVertex[2] (Sink)
  ├── ExecutionVertex[2][0]  →  Execution(attempt=0, state=CREATED)
  ├── ... (4 个)

三层结构

层级对应数量
1ExecutionJobVertexJobVertex1 per operator group
2ExecutionVertexSubtaskparallelism 个
3ExecutionAttempt多个(每次重试 +1)

4. DefaultScheduler—调度核心

4.1 调度流程

// DefaultScheduler.java — 调度入口
public class DefaultScheduler extends SchedulerBase {
    private final SchedulingStrategy schedulingStrategy;   // Eager / Lazy-From-Sources
    private final ExecutionSlotAllocator slotAllocator;    // Slot 分配器
    private final ExecutionDeployer executionDeployer;     // Task 部署器
    private final ExecutionFailureHandler failureHandler;  // 失败处理器

    @Override
    protected void startSchedulingInternal() {
        // 1. 转换状态: CREATED → RUNNING
        transitionToRunning();  // executionGraph.transitionToRunning()

        // 2. 调度策略开始调度
        schedulingStrategy.startScheduling();
        // → EagerSchedulingStrategy: 所有 Vertex 同时分配 Slot
        // → PipelinedRegionSchedulingStrategy: 从 Source 开始逐步分配
    }

    // 为指定 Vertex 分配 Slot 并部署
    public void allocateSlotsAndDeploy(List<ExecutionVertexID> vertices) {
        // 1. 获取每个 Vertex 的当前 Execution
        List<Execution> executions = vertices.stream()
            .map(vid -> executionGraph.getExecutionVertex(vid).getCurrentExecution())
            .collect(toList());

        // 2. 委托给 ExecutionDeployer
        executionDeployer.allocateSlotsAndDeploy(executions);
        // 3. → ExecutionSlotAllocator.allocateSlotsFor(executions)
        // 4. → SlotPoolService.requestNewAllocatedSlot() / allocateAvailableSlot()
        // 5. → Slot 分配成功 → Execution.deploy()
    }
}

4.2 调度策略

策略行为适用
Eager所有 Vertex 同时分配 Slot资源充足、小 Job
Lazy-From-Sources从 Source 开始,按拓扑顺序逐步分配节省 Slot、大 Job
PipelinedRegion以 Pipelined Region 为单位调度(默认)批流统一

4.3 SlotPool — JobMaster 端 Slot 管理

// DeclarativeSlotPoolBridge.java — SlotPool 的核心实现
public class DeclarativeSlotPoolBridge implements SlotPoolService {
    // 待处理的 Slot 请求
    private final Map<SlotRequestId, PendingRequest> pendingRequests;

    // 资源需求声明 (发送给 RM)
    private final DeclarativeSlotPool declarativeSlotPool;

    // 请求新 Slot
    public CompletableFuture<PhysicalSlot> requestNewAllocatedSlot(
            SlotRequestId requestId, ResourceProfile profile) {

        // 1. 创建 PendingRequest (带超时)
        PendingRequest request = new PendingRequest(requestId, profile, timeout);
        pendingRequests.put(requestId, request);

        // 2. 增加资源需求
        declarativeSlotPool.increaseResourceRequirementsBy(profile, 1);
        // → 通知 SlotPoolService → RM.declareRequiredResources()

        return request.getFuture();
    }

    // TM 提供 Slot 后的回调
    public void newSlotsAreAvailable(Collection<PhysicalSlot> slots) {
        // 1. 匹配: 将 Slot 分配给 PendingRequest
        Map<SlotRequestId, PhysicalSlot> matches =
            requestSlotMatchingStrategy.matchRequestsAndSlots(pendingRequests, slots);

        // 2. 满足匹配的请求
        for (var match : matches.entrySet()) {
            PendingRequest request = pendingRequests.remove(match.getKey());
            request.complete(match.getValue());  // → ExecutionDeployer.deploy()
        }
    }
}

5. TaskManager 端—Task 生命周期

5.1 Task.doRun() — 核心执行

// Task.java — Task 执行流程
public void doRun() {
    // 1. 状态转换: CREATED → DEPLOYING
    executionState = ExecutionState.DEPLOYING;

    try {
        // 2. Bootstrap: 加载类、建立网络连接、创建 Environment
        //    - 从 BlobServer 下载 JAR
        //    - 创建 ResultPartition + InputGate
        //    - 加载 TaskInvokable 类 (如 StreamTask)
        //    - 创建 RuntimeEnvironment (Task 的运行时上下文)

        // 3. 恢复和调用
        restoreAndInvoke(invokable);
        // ↓

    } catch (Throwable t) {
        // 4. 错误处理: → FAILED / CANCELED
    } finally {
        // 5. 清理: 释放网络、内存、文件系统
    }
}

private void restoreAndInvoke(TaskInvokable invokable) throws Exception {
    // 3a. DEPLOYING → INITIALIZING
    setExecutionState(ExecutionState.INITIALIZING, ...);
    invokable.restore();  // StateBackend.restore()

    // 3b. INITIALIZING → RUNNING
    setExecutionState(ExecutionState.RUNNING, ...);
    invokable.invoke();   // 进入数据处理循环
}

5.2 Task Execution State 状态机

CREATED ──► DEPLOYING ──► INITIALIZING ──► RUNNING ──► FINISHED
                          │                │
                          │   (cancel/fail)│
                          ▼                ▼
                      CANCELED/CANCELING  FAILED

6. ExecutionGraph 的 JobStatus 状态机

CREATED ──[startSchedulingInternal]──► RUNNING
                                          │
                    ┌─────────────────────┤
                    ▼                     ▼
              RESTARTING             FINISHED
              (失败后可重试)          (所有 Task 完成)
                    │
                    ▼
                RUNNING
                    │
          ┌─────────┴─────────┐
          ▼                   ▼
      CANCELLING           FAILING
          │                   │
          ▼                   ▼
      CANCELED              FAILED
// ExecutionGraph.java — 状态转换
public void transitionState(JobStatus targetStatus) {
    synchronized (lock) {
        JobStatus current = state;
        if (current.isTerminalState()) return;  // 终端状态不可逆

        // 状态一致性校验
        switch (targetStatus) {
            case RUNNING:     assert current == CREATED || current == RESTARTING; break;
            case FINISHED:    assert current == RUNNING; break;
            case RESTARTING:  assert current == RUNNING; break;
            case CANCELLING:  assert current == RUNNING || current == RESTARTING; break;
            case FAILING:     assert current == RUNNING || current == RESTARTING; break;
        }

        state = targetStatus;
        stateTimestamps[targetStatus.ordinal()] = System.currentTimeMillis();
        notifyJobStatusListeners(current, targetStatus);
    }
}

7. Failover — 容错恢复

7.1 恢复决策

// DefaultScheduler.java — 失败处理
public void handleTaskFailure(Execution execution, Throwable error) {
    // 1. 调用失败策略
    FailureHandlingResult result = executionFailureHandler.getFailureHandlingResult(
        execution.getVertex().getID(), error);

    // 2. 决定是否重启
    if (result.canRestart()) {
        restartTasksWithDelay(result);  // 延迟后重启
    } else {
        failJob(error);                  // 不可重启 → Job 失败
    }
}

// ExecutionFailureHandler — 结合多种策略
public FailureHandlingResult getFailureHandlingResult(
        ExecutionVertexID failedVertex, Throwable error) {

    // 1. FailoverStrategy: 决定哪些 Vertex 需要重启
    //    - RestartAllStrategy: 所有 Vertex 都重启
    //    - RestartPipelinedRegionStrategy: 只重启受影响的 Region
    Set<ExecutionVertexID> verticesToRestart =
        failoverStrategy.getTasksNeedingRestart(failedVertex, error);

    // 2. RestartBackoffTimeStrategy: 决定是否及何时重启
    //    - NoRestart: 不重启
    //    - FixedDelay: 最多 N 次,每次等 delay
    //    - FailureRate: 时间窗口内最多 N 次
    //    - ExponentialDelay: 指数退避
    if (backoffStrategy.canRestart()) {
        long restartDelay = backoffStrategy.getBackoffTime();
        return FailureHandlingResult.restartable(verticesToRestart, restartDelay);
    }

    return FailureHandlingResult.unrecoverable(error);
}

7.2 重启流程

// DefaultScheduler.java — 重启流程
private void restartTasks(Set<ExecutionVertexVersion> verticesToRestart,
                          boolean isGlobalRecovery) {

    for (ExecutionVertexVersion version : verticesToRestart) {
        ExecutionVertex vertex = executionGraph.getExecutionVertex(version.getVertexId());

        // 1. 重置 ExecutionVertex: 归档旧 Execution, 创建新 Execution
        vertex.resetForNewExecution();
        // → currentExecution = new Execution(vertex, attemptNumber++)
        // → 旧 Execution 归档到 executionHistory

        // 2. 恢复状态
        if (isGlobalRecovery) {
            // 全局恢复: 从最近 Checkpoint 恢复所有状态
            vertex.restoreState(latestCheckpointState);
        } else {
            // 局部恢复: 只恢复受影响的 Subtask 状态
            vertex.restoreSubtaskState(subtaskStateHandle);
        }
    }

    // 3. 重新部署
    schedulingStrategy.restartTasks(verticesToRestart);
    // → allocateSlotsAndDeploy(restartedVertices)
}

8. 端到端调度全链路总结

1 Client:          CliFrontend.run() → StreamGraph → JobGraph
2 Dispatcher:      submitJob(JobGraph) → persistAndRunJob()
3 JobMaster:       createJobMasterRunner → JobMaster.onStart()
4 JobMaster:       startScheduling() → schedulerNG.startScheduling()
5 DefaultScheduler: startSchedulingInternal() → transitionToRunning()
6 SchedulingStrategy: startScheduling() → 挑选就绪的 Vertex
7 DefaultScheduler: allocateSlotsAndDeploy(vertices)
8 ExecutionDeployer: → Execution.transitionState(SCHEDULED)
9 SlotPool:          → requestNewAllocatedSlot(resourceProfile)
10 ResourceManager:   → allocateSlot() → 在 TM 上分配 Slot
11 TM:                → offerSlots → newSlotsAreAvailable()
12 ExecutionDeployer: → Execution.deploy()
13 RPC:               → TaskExecutor.submitTask(TaskDeploymentDescriptor)
14 TM:                → Task.doRun(): CREATED → DEPLOYING → INITIALIZING → RUNNING
15 StreamTask:        → invoke() → processInput() 循环

On Failure:
16 TM:                → Task → FAILED
17 JM:                → updateTaskExecutionState()
18 DefaultScheduler:  → handleTaskFailure()
19 ExecutionFailureHandler: → canRestart? restartDelay?
20 DefaultScheduler:  → restartTasks()
21 ExecutionVertex:   → resetForNewExecution() → 新 attempt
    → 回到步骤 7

9. 源码导航(完整版)

文件关键类/方法职责
flink-clients/.../CliFrontend.javamain(), run()CLI 入口
flink-streaming-java/.../StreamGraphGenerator.javagenerate()DataStream → StreamGraph
flink-optimizer/.../StreamingJobGraphGenerator.javacreateJobGraph(), isChainable()StreamGraph → JobGraph
flink-runtime/.../dispatcher/Dispatcher.javasubmitJob(), persistAndRunJob(), onStart()Job 提交与恢复
flink-runtime/.../jobmaster/JobMaster.javaonStart(), startScheduling()Per-Job Master
flink-runtime/.../scheduler/DefaultScheduler.javastartSchedulingInternal(), allocateSlotsAndDeploy(), restartTasks()调度核心
flink-runtime/.../scheduler/SchedulerBase.javaupdateTaskExecutionState(), startScheduling()调度基类
flink-runtime/.../executiongraph/ExecutionGraph.javascheduleForExecution(), transitionState(), restart()物理执行图
flink-runtime/.../executiongraph/ExecutionVertex.javaresetForNewExecution(), deploy()单 Subtask
flink-runtime/.../executiongraph/Execution.javadeploy(), transitionState()单次执行尝试
flink-runtime/.../jobmaster/slotpool/DeclarativeSlotPoolBridge.javarequestNewAllocatedSlot(), newSlotsAreAvailable()Slot Pool
flink-runtime/.../resourcemanager/ResourceManager.javaregisterTaskManager(), declareRequiredResources()RM 基类
flink-runtime/.../taskmanager/Task.javadoRun(), restoreAndInvoke()Task 生命周期

10. 常见问题 / 面试题

Q1: JobGraph 和 ExecutionGraph 的本质区别?

A: JobGraph 是 Client 端生成的逻辑图(JobVertex + IntermediateDataSet),不包含并行度和部署信息。ExecutionGraph 是 JobMaster 端生成的物理图,包含并行度(ExecutionVertex)、Slot 分配、Execution 状态。JobGraph 可以被持久化和跨 Flink 版本重用,ExecutionGraph 是瞬态的。

Q2: 为什么 Flink 的 Failover 比 Spark 快?

A: Flink (1) Checkpoint 粒度是算子级,每个 Task 独立做快照 → 恢复时只需恢复失败的 Task(Region Recovery);(2) 不需要重新计算上游数据(Source offset 从 Checkpoint 回滚);(3) Slots 已预热(TaskManager 仍在运行)。Spark 的 Stage 级恢复需要从 RDD 源头重新计算整个 Stage。

Q3: Dispatcher 的去重机制如何工作?

A: 两层去重:(1) JobResultStore 检查全局终端状态(FINISHED/FAILED/CANCELED)→ 如果已存在 → 幂等返回;(2) submittedAndWaitingTerminationJobIDs 检查本地是否正在运行/等待终止 → 防止并发重复提交。这确保了 flink run 命令的重试幂等性。

Q4: SlotPool 和 ResourceManager 的分工是什么?

A: SlotPool (JobMaster 端) 维护 Job 需要的资源需求声明,管理从 TM 获得的 Slot。ResourceManager (集群端) 管理全局的 Slot(所有 TM 的 Slot),负责分配 Slot 给各 Job 的 SlotPool。SlotPool 只关心自己 Job 需要多少资源,RM 关心整个集群的资源分配。

Q5: DefaultScheduler 的 Lazy-From-Sources 策略何时生效?

A: 当 slot.request.timeout 配置值 ≠ 0 时(非 Eager 模式)。Scheduler 先只分配 Source 的 Slot,Source 部署完成后产生数据 → 触发下游的 IntermediateResultPartition 可用 → 下游 Vertex 的 InputGate ready → 调度下游。好处是逐步分配 Slot,避免空闲 Slot 占用。

Q6: Execution 的 attempt 号是做什么的?

A: 每次失败重启时 attemptNumber++,用于:(1) 区分同一 Vertex 的不同尝试 → Checkpoint 文件中 (checkpointId, attemptNumber) 唯一标识;(2) 防止 stale 消息(late Ack from old attempt);(3) 清除上次 attempt 的 Shuffle 数据(通过 attemptNumber 命名 ResultPartition)。


下一步