The source notes for this track are currently maintained in Chinese.
预计阅读时间: 120 分钟 前置阅读: doc-00——架构概览与代码地图 下一次阅读: doc-02(流数据链路), doc-11(ResourceManager)

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
| 维度 | StreamGraph | JobGraph |
|---|---|---|
| 粒度 | 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 个)
三层结构:
| 层级 | 类 | 对应 | 数量 |
|---|---|---|---|
| 1 | ExecutionJobVertex | JobVertex | 1 per operator group |
| 2 | ExecutionVertex | Subtask | parallelism 个 |
| 3 | Execution | Attempt | 多个(每次重试 +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.java | main(), run() | CLI 入口 |
flink-streaming-java/.../StreamGraphGenerator.java | generate() | DataStream → StreamGraph |
flink-optimizer/.../StreamingJobGraphGenerator.java | createJobGraph(), isChainable() | StreamGraph → JobGraph |
flink-runtime/.../dispatcher/Dispatcher.java | submitJob(), persistAndRunJob(), onStart() | Job 提交与恢复 |
flink-runtime/.../jobmaster/JobMaster.java | onStart(), startScheduling() | Per-Job Master |
flink-runtime/.../scheduler/DefaultScheduler.java | startSchedulingInternal(), allocateSlotsAndDeploy(), restartTasks() | 调度核心 |
flink-runtime/.../scheduler/SchedulerBase.java | updateTaskExecutionState(), startScheduling() | 调度基类 |
flink-runtime/.../executiongraph/ExecutionGraph.java | scheduleForExecution(), transitionState(), restart() | 物理执行图 |
flink-runtime/.../executiongraph/ExecutionVertex.java | resetForNewExecution(), deploy() | 单 Subtask |
flink-runtime/.../executiongraph/Execution.java | deploy(), transitionState() | 单次执行尝试 |
flink-runtime/.../jobmaster/slotpool/DeclarativeSlotPoolBridge.java | requestNewAllocatedSlot(), newSlotsAreAvailable() | Slot Pool |
flink-runtime/.../resourcemanager/ResourceManager.java | registerTaskManager(), declareRequiredResources() | RM 基类 |
flink-runtime/.../taskmanager/Task.java | doRun(), 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)。
下一步
- 流数据链路: doc-02-stream-dataflow.md——Source → Transform → Sink 的端到端数据流
- ResourceManager 深度: doc-11-resource-scheduling.md——Slot 分配与管理