Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 01

Job submission and scheduling

Trace how a JAR or SQL job becomes distributed tasks through Flink's submission and scheduling path.

Reading
120 min
Track
Flink internals
Source
Chinese source notes

The source notes for this track are currently maintained in Chinese.

预计阅读时间: 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)。


下一步