ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

第 03 篇:「作业图与调度器」—— StreamGraph → JobGraph → ExecutionGraph 三级转换与 DefaultScheduler 调度内核

2026/9/5 6:13:18 拓冰建站 浏览量
第 03 篇:「作业图与调度器」—— StreamGraph → JobGraph → ExecutionGraph 三级转换与 DefaultScheduler 调度内核 仓库https://github.com/apache/flink官方文档https://nightlies.apache.org/flink/flink-docs-lts/技术栈Java 11 / StreamGraph / JobGraph / ExecutionGraph / Scheduler解读版本release-1.20.5commit0980485解读视角总架构师评审架构 / 源码 / 生产 / 进阶第 03 篇「作业图与调度器」—— StreamGraph → JobGraph → ExecutionGraph 三级转换与 DefaultScheduler 调度内核阅读本文你将了解一条 DataStream/SQL 程序要经历StreamGraph → JobGraph → ExecutionGraph三级图转换才能变成可分布式执行的 Task对应internals/job_scheduling.md。算子链Operator Chaining在StreamingJobGraphGenerator.createChain里递归合并判据收敛在isChainable——同并行度 ForwardPartitioner ChainingStrategy 允许flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:1543。JobGraph → ExecutionGraph的转换入口是DefaultExecutionGraphBuilder它把每个JobVertex按并行度展开成ExecutionVertex再挂到ExecutionGraph.attachJobGraphflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/DefaultExecutionGraphBuilder.java:225。调度的真正起点是SchedulerBase.startScheduling它下钻到DefaultScheduler.startSchedulingInternal再交给PipelinedRegionSchedulingStrategy按 Pipelined Region 拓扑排序逐个调度flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:234。调度策略最终回落到allocateSlotsAndDeploy申请 Slot 部署 Task是调度器的原子动作flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:466。03.0 一句话定性Flink 的作业调度可以定性为用三级图模型StreamGraph → JobGraph → ExecutionGraph把用户逻辑逐步翻译成可调度的物理执行计划再用一个基于 Pipelined Region 的调度策略按拓扑顺序把执行顶点映射到 Slot 上。StreamGraph 是离用户最近的图算子、连接、并行度都在但还带着 DSL 痕迹JobGraph 是提交给集群的图算子链已合并、任务顶点已确定ExecutionGraph 是运行期真正被调度的图每个 JobVertex 展开成并行度个 ExecutionVertex每个 ExecutionVertex 又可能有多轮 Execution 尝试。这三级图的关系本质是编译的三次 pass。03.1 架构视角Architect03.1.1 三级图模型谁负责什么图生成时机粒度关键结构源码位置StreamGraph客户端StreamExecutionEnvironment构建算子 连接StreamNode/StreamEdgeflink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamGraph.java:89JobGraph客户端提交前任务顶点算子链已合并JobVertex/IntermediateDataSetflink-runtime/src/main/java/org/apache/flink/runtime/jobgraph/JobGraph.java:135ExecutionGraphJobMaster 运行期并行子任务 执行尝试ExecutionVertex/Executionflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:79StreamGraph的节点是StreamNode、边是StreamEdge它由StreamGraphGenerator.generate()flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamGraphGenerator.java:274从 DataStream API 的transform调用链生成。它离用户最近还保留着每个算子一个节点的原始粒度。JobGraph是给集群看的图节点是JobVertex每个JobVertex已经是一段算子链多个可 chaining 的算子被合并。它只通过addVertex注册顶点flink-runtime/src/main/java/org/apache/flink/runtime/jobgraph/JobGraph.java:291setSnapshotSettings挂上 checkpoint 配置:359。ExecutionGraph是运行期真正被调度的图。它的三层结构在ExecutionGraph接口的 Javadoc 里写得很清楚flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:56-77ExecutionJobVertex对应 JobGraph 里的一个JobVertex通常是一个 operatorExecutionVertex对应一个并行子任务Execution是一次执行尝试失败重试会生成新的Execution用ExecutionAttemptID标识。03.1.2 三级图转换的状态图下面这张状态图把一个程序从代码到可调度计划的三次转换串起来图示讲解这张状态图回答程序怎么从代码一步步变成调度器眼里的拓扑。第一步StreamGraphGenerator.generate()flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamGraphGenerator.java:274把 DataStream 的算子 DAG 固化成StreamGraph此时每个算子一个StreamNode第二步StreamingJobGraphGenerator.createJobGraph()flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:242做算子链合并产出JobGraph第三步DefaultExecutionGraphBuilder.buildGraph()flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/DefaultExecutionGraphBuilder.java:225在 JobMaster 侧把JobGraph展开成ExecutionGraph并按并行度实例化ExecutionVertex第四步ExecutionGraph.attachJobGraph()flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:158完成顶点连接与拓扑排序产出SchedulingTopology第五步DefaultScheduler从这个拓扑开始调度flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:234。注意前两步发生在客户端后三步发生在 JobMaster——这是理解哪些优化离线做、哪些决策在线做的关键分界。03.2 源码侦探Source Sleuth03.2.1 StreamGraph → JobGraphcreateJobGraph 的编排顺序StreamingJobGraphGenerator.createJobGraph()是三级转换里最关键的一跳它把StreamGraph重写为JobGraph。私有方法createJobGraph()flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:242的执行顺序本身就是一份转换 checklist// flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:242privateJobGraphcreateJobGraph(){preValidate();jobGraph.setJobType(streamGraph.getJobType());jobGraph.setDynamic(streamGraph.isDynamic());// 1. 为每个节点生成确定性哈希跨提交识别同一算子用于状态恢复MapInteger,byte[]hashesdefaultStreamGraphHasher.traverseStreamGraphAndGenerateHashes(streamGraph);// 2. 算子链合并核心步骤setChaining(hashes,legacyHashes);// 3. 物理边、SlotSharingGroup、内存占比、checkpoint 配置……setPhysicalEdges();setSlotSharingAndCoLocation();configureCheckpointing();jobGraph.setSavepointRestoreSettings(streamGraph.getSavepointRestoreSettings());// ...}关键源码事实createJobGraph()的编排顺序揭示了 Flink 的一条设计原则——先算哈希、再合并算子链flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:252-261。算子哈希traverseStreamGraphAndGenerateHashes必须在 chaining 之前算因为链合并后多个算子共享一个JobVertex它们的哈希要用来生成OperatorID保证同一个算子在下一次提交时仍被认出来——这是状态恢复savepoint 按 OperatorID 匹配状态的前提。顺序错了算子 ID 就会漂移savepoint 直接失效。03.2.2 算子链Operator Chaining的合并内核算子链是 Flink 性能的基石把不需要网络交换的相邻算子合并进同一个 Task省掉序列化/反序列化与网络往返。合并的递归入口是createChain// flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:670privateListStreamEdgecreateChain(finalIntegercurrentNodeId,finalintchainIndex,finalOperatorChainInfochainInfo,finalMapInteger,OperatorChainInfochainEntryPoints){// ...for(StreamEdgeoutEdge:currentNode.getOutEdges()){if(isChainable(outEdge,streamGraph)){chainableOutputs.add(outEdge);// 可合并递归下钻}else{nonChainableOutputs.add(outEdge);// 不可合并断开新起一条链}}for(StreamEdgechainable:chainableOutputs){transitiveOutEdges.addAll(createChain(chainable.getTargetId(),chainIndex1,chainInfo,chainEntryPoints));}// ...}关键源码事实createChain用递归把一条链上的算子逐个收编flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:670-725。它对当前节点的每条出边调用isChainable判定可合并的出边chainIndex 1递归下钻不可合并的出边chainEntryPoints.computeIfAbsent新建一条链的入口。递归的边界就是遇到不可合并的边——那一刻前面积累的算子就固化成同一个JobVertex。chainIndex从 1 开始0 预留给 chained source input这个编号最终决定算子在OperatorChain内部的位置。判据本身收敛在isChainable与isChainableInput// flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:1543publicstaticbooleanisChainable(StreamEdgeedge,StreamGraphstreamGraph){StreamNodedownStreamVertexstreamGraph.getTargetVertex(edge);returndownStreamVertex.getInEdges().size()1isChainableInput(edge,streamGraph);}isChainableInput里的硬条件flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamingJobGraphGenerator.java:1549-1570有三条①isChainingEnabled()全局开关 上下游同SlotSharingGroupareOperatorsChainablearePartitionerAndExchangeModeChainable② 下游不能是 union 的多输入inEdge.getTypeNumber()相同则不可合并③ 并行度必须一致upStreamVertex.getParallelism() downStreamVertex.getParallelism()在:1642。这三条共同决定了什么时候必须跨网络、什么时候可以塞进一个 Task。03.2.3 JobGraph → ExecutionGraphattachJobGraph 的展开JobGraph到了 JobMaster 后由DefaultExecutionGraphBuilder展开成ExecutionGraph。真正把JobGraph挂到ExecutionGraph的动作是attachJobGraph// flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:158voidattachJobGraph(ListJobVertextopologicallySorted,JobManagerJobMetricGroupjobManagerJobMetricGroup)throwsJobException;关键源码事实attachJobGraph的入参是拓扑排序后的 JobVertex 列表flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:158。拓扑排序在DefaultExecutionGraphBuilder里先做掉flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/DefaultExecutionGraphBuilder.java:225保证挂接时顶点按数据流方向有序。挂接过程中每个JobVertex被包装成ExecutionJobVertex再按并行度initializeJobVertexflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:240实例化出parallelism个ExecutionVertex。这一步是逻辑任务到物理子任务的分水岭——JobGraph里一个JobVertex并行度是 32到这里就变成 32 个可独立调度、独立失败的ExecutionVertex。03.2.4 调度入口SchedulerBase.startScheduling 的下钻调度器的统一入口是SchedulerBase.startScheduling它是一个final方法锁死调度骨架把可变部分留给子类的startSchedulingInternal// flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/SchedulerBase.java:625publicfinalvoidstartScheduling(){mainThreadExecutor.assertRunningInMainThread();registerJobMetrics(/* ... */);operatorCoordinatorHandler.startAllOperatorCoordinators();startSchedulingInternal();}DefaultScheduler的实现只有三步flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:234打日志 →transitionToRunning()→schedulingStrategy.startScheduling()。transitionToRunning最终调到executionGraph.transitionToRunning()flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/SchedulerBase.java:564把作业状态从CREATED切到RUNNING随后把调度权交给策略。03.2.5 调度策略PipelinedRegionSchedulingStrategy 的拓扑推进默认的调度策略是PipelinedRegionSchedulingStrategy它按Pipelined Region一段完全由 pipelined 边连通的子图为粒度调度而非逐个顶点调度// flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/strategy/PipelinedRegionSchedulingStrategy.java:282privatevoidscheduleRegion(finalSchedulingPipelinedRegionregion){checkState(areRegionVerticesAllInCreatedState(region),BUG: trying to schedule a region which is not in CREATED state);scheduledRegions.add(region);schedulerOperations.allocateSlotsAndDeploy(regionVerticesSorted.get(region));}关键源码事实scheduleRegion把调度一个 region收敛为一次allocateSlotsAndDeploy调用flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/strategy/PipelinedRegionSchedulingStrategy.java:282-288。调度策略自己不碰 Slot 申请只负责决定哪些顶点该被调度、按什么顺序把申请 Slot 部署这个动作通过SchedulerOperations回抛给DefaultScheduler。startScheduling()从源 regionisSourceRegion判定:182开始maybeScheduleRegions里用SchedulingStrategyUtils.sortPipelinedRegionsInTopologicalOrder:231按拓扑序依次调度保证上游 region 的 blocking 边数据产完下游 region 才被放行。这个设计把调度顺序和资源分配彻底解耦。allocateSlotsAndDeploy是调度器与 Slot 管理器打交道的原子动作// flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:466publicvoidallocateSlotsAndDeploy(finalListExecutionVertexIDverticesToDeploy){finalMapExecutionVertexID,ExecutionVertexVersionrequiredVersionByVertexexecutionVertexVersioner.recordVertexModifications(verticesToDeploy);finalListExecutionexecutionsToDeployverticesToDeploy.stream().map(this::getCurrentExecutionOfVertex).collect(Collectors.toList());executionDeployer.allocateSlotsAndDeploy(executionsToDeploy,requiredVersionByVertex);}关键源码事实allocateSlotsAndDeploy先recordVertexModifications记录版本号再委托executionDeployer.allocateSlotsAndDeploy去真正申请 Slotflink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:466-476。这个ExecutionVertexVersion版本号是并发安全的守卫recordVertexModifications会给这批顶点打上新版本号部署动作完成后若版本号已变说明期间发生了 failover 重置就能识别出过期部署并丢弃。这是调度器处理调度进行中任务失败这一竞态的基石。03.2.6 调度时序图从 submitJob 到 deploy下面这张时序图把作业提交 → 图转换 → 调度 → 部署完整链路串起来图示讲解这张时序图回答一个 JobGraph 如何一步步变成跑在 TaskExecutor 上的 Task。链路分五段第一段客户端把JobGraph交给Dispatcher.submitJobflink-runtime/src/main/java/org/apache/flink/runtime/dispatcher/Dispatcher.java:518Dispatcher创建JobMaster第二段JobMaster用DefaultExecutionGraphBuilder.buildGraph()flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/DefaultExecutionGraphBuilder.java:225把JobGraph展开成ExecutionGraphattachJobGraphflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:158实例化出ExecutionVertex第三段JobMaster.startScheduling()flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMaster.java:1235调用schedulerNG.startScheduling()进入DefaultScheduler.startSchedulingInternalflink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:234先transitionToRunning:238再schedulingStrategy.startScheduling():239第四段策略按 Pipelined Region 拓扑序推进scheduleRegion回调allocateSlotsAndDeployflink-runtime/src/main/java/org/apache/flink/runtime/scheduler/strategy/PipelinedRegionSchedulingStrategy.java:282向ResourceManager申请 Slot第五段Slot 就绪后调度器经Execution.deployflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java:557把 Task 部署到TaskExecutorTask.doRun()flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:581启动物理执行。注意调度策略与 Slot 申请是两条腿——策略只决定顺序资源分配走独立的ResourceManager通道这正是 Flink 调度可插拔性的体现。03.3 生产实践Production03.3.1 算子链与 SlotSharingGroup 的配置权衡配置项作用生产建议pipeline.operator-chaining全局算子链开关默认开启关掉会退化为每算子一 Task吞吐暴跌disableChaining()/startNewChain()算子级关闭/断开链定位单个算子延迟时临时用slotSharingGroup(name)强制同组算子共享 Slot默认所有算子一个 default 组pipeline.max-parallelismmaxParallelismkey group 数与 02 篇联动上线前规划算子链不是越多越好链太长会把一条链上的背压放大——链尾一个算子慢整条链的输入 buffer 都被拖住。生产上默认全链、按需断开是主流策略用 Flink UI 的 Task 明细定位到具体算子后对该算子disableChaining()再发布把它单独拆出来观察是最快的定位手段。03.3.2 调度相关的三个高频问题Slot 不够导致作业一直 PENDING调度器在allocateSlotsAndDeploy后等待 Slot 就绪若ResourceManager无法满足资源作业停在SCHEDULED/DEPLOYING不前进。排查入口是 TaskManager 的requestSlot日志flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:1183与 RM 的 slot 分配计数。blocking 边导致下游迟迟不启动批作业里PipelinedRegionSchedulingStrategy会等上游 region 的 blocking 边数据ALL_DATA_PRODUCED才放行下游flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/strategy/PipelinedRegionSchedulingStrategy.java:327这是调度慢的常见误判——其实是在等数据落盘。failover 引发版本号竞态recordVertexModifications的版本号机制flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/DefaultScheduler.java:468会丢弃过期部署若观察到任务被重复部署又立刻取消往往是因为 failover 与正常调度并发属正常防御而非 bug。03.4 进阶向导Advanced03.4.1 调度器族谱Default vs Adaptive1.20 的调度器已经家族化DefaultScheduler只是默认成员调度器适用场景关键差异DefaultScheduler流作业 固定并行度批作业并行度固定按 Pipelined Region 调度AdaptiveScheduler批作业、资源动态变化运行时根据可用 Slot 动态调整并行度AdaptiveBatchScheduler批作业、大状态自适应 优化 shuffle 落盘DefaultScheduler与AdaptiveScheduler都继承SchedulerBase共享startScheduling骨架flink-runtime/src/main/java/org/apache/flink/runtime/scheduler/SchedulerBase.java:625差异只在startSchedulingInternal与调度策略工厂。这个骨架锁死、策略可插拔的设计正是 Flink 能在不重写调度器的情况下演进出自适应调度的原因。03.4.2 与 Spark 的 DAG 调度对比维度FlinkSpark图模型StreamGraph→JobGraph→ExecutionGraph 三级RDD DAG→物理 Stage 两级调度粒度Pipelined Region流 / 逐个 ExecutionVertexStage宽依赖切分算子合并Operator Chaining同并行度Forward同一 Stage 内 pipeline 执行资源模型SlotTaskManager 内固定分片Executor Core/内存失败恢复基于 checkpoint 的细粒度/全局恢复血缘重算RDD lineage03.4.3 一个容易被误读的点很多人以为StreamGraph是用户代码里就能拿到的东西。实际上StreamGraph是StreamExecutionEnvironment内部的中间产物用户代码触发execute()时才由StreamGraphGenerator.generate()惰性生成flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamGraphGenerator.java:274随后立刻被StreamingJobGraphGenerator转成JobGraph。用户能直接拿到的最上层图其实是StreamExecutionEnvironment.getExecutionPlan()返回的 JSON 计划对应ExecutionGraph的可视化版本。这个认知差异在想在客户端拦截图做自定义优化时很关键——真正的钩子点不是 StreamGraph而是StreamGraphGenerator的transform链或JobGraph提交前。03.5 小结与下一篇本篇打穿了 Flink 调度机制的一条主干三级图转换StreamingJobGraphGenerator.createJobGraph的算子链合并 →DefaultExecutionGraphBuilder的并行度展开与DefaultScheduler 的拓扑调度PipelinedRegionSchedulingStrategy按 region 推进 →allocateSlotsAndDeploy申请 Slot 部署 Task。理解了调度顺序策略与资源分配Slot的解耦就理解了 Flink 调度器能插拔、能自适应的根源。下一篇进入《Task 生命周期与 TaskManager》切入点是调度器部署出 Task 之后的故事Task 的CREATED→DEPLOYING→INITIALIZING→RUNNING→FINISHED状态机如何用 CAS 自旋处理部署竞态以及 TaskExecutor 的 Slot 管理与心跳如何让 JobManager 知道这个 Task 还活着对应本系列 04 篇。