第 04 篇:「Task 生命周期与 TaskManager」—— Task 状态机 CAS 自旋与 TaskExecutor 的 Slot/心跳 📅 发布时间:2026/9/5 5:50:53 👁 浏览次数: 仓库https://github.com/apache/flink官方文档https://nightlies.apache.org/flink/flink-docs-lts/技术栈Java 11 / Task / ExecutionState / TaskExecutor / Slot / Heartbeat解读版本release-1.20.5commit0980485解读视角总架构师评审架构 / 源码 / 生产 / 进阶第 04 篇「Task 生命周期与 TaskManager」—— Task 状态机 CAS 自旋与 TaskExecutor 的 Slot/心跳阅读本文你将了解Task 的物理执行走CREATED → DEPLOYING → INITIALIZING → RUNNING → FINISHED状态机由Task.doRun的while(true)自旋 CAS 驱动对应internals/task_lifecycle.md。状态迁移的核心是AtomicReferenceFieldUpdater.compareAndSet一次失败就重试这是部署竞态的源码级解法flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:1100。Task 每进入INITIALIZING/RUNNING都通过taskManagerActions.updateTaskExecutionState向 JobManager 上报最终状态由notifyFinalState收敛flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:1069。TaskExecutor 启动时是反向寻址——先resourceManagerLeaderRetriever.start去找 RM而非等 RM 来连它flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:487。心跳分两路对 JobManager 的heartbeatFromJobManager与对 ResourceManager 的heartbeatFromResourceManager各自独立超时flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:1063、:1069。04.0 一句话定性Flink 的 Task 执行可以定性为用AtomicReferenceFieldUpdater的 CAS 自旋把一个物理 Task 状态机跑在 TaskExecutor 的 Slot 里同时用两路独立心跳对 JM、对 RM维持我还活着的分布式认知。这里要区分两个常常被混为一谈的状态JobMaster 侧的Execution逻辑执行尝试状态叫ExecutionState与 TaskManager 侧的Task物理执行线程也有自己的状态机。两者通过updateTaskExecutionStateRPC 同步构成逻辑调度与物理执行的双层状态模型。04.1 架构视角Architect04.1.1 ExecutionState官方状态机的全貌ExecutionState枚举的 Javadoc 里自带一张 ASCII 状态图flink-runtime/src/main/java/org/apache/flink/runtime/execution/ExecutionState.java:21-47是理解 Task 生命周期最权威的来源CREATED - SCHEDULED - DEPLOYING - INITIALIZING - RUNNING - FINISHED | | | | | | | | ------------------- | | V V | | CANCELLING --------- CANCELED | | | | ------------------------- | | ... - FAILED V RECONCILING - INITIALIZING | RUNNING | FINISHED | CANCELED | FAILED主线是CREATED → SCHEDULED → DEPLOYING → INITIALIZING → RUNNING → FINISHED任何状态都能进入FAILEDFINISHED/CANCELED/FAILED是三个终态isTerminal()判定的正是这三个flink-runtime/src/main/java/org/apache/flink/runtime/execution/ExecutionState.java:76。RECONCILING是 JobManager failover 后、Task 与 JM 状态对账的中间态。04.1.2 双层状态模型Execution逻辑与 Task物理层类生命周期谁持有状态枚举逻辑Execution一次执行尝试可重试JobMasterExecutionGraphExecutionState物理Task一个真实线程跑用户代码TaskExecutorSlot私有ExecutionState字段Execution是 JobMaster 眼里的这个 ExecutionVertex 的第 N 次尝试用ExecutionAttemptID标识它的deploy()flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java:557负责把任务打包发给 TaskExecutor。Task是 TaskExecutor 上真正run()的线程flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:573。两者状态靠updateTaskExecutionState消息同步——Task 每变一次状态就回传一次。04.1.3 TaskExecutor一个进程三块职责TaskExecutorflink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:203是 TaskManager 的核心RpcEndpoint职责分三块Slot 管理TaskSlotTable见:354、心跳维护对 JM 与 RM 两个HeartbeatManager见:381/:383、任务执行接受submitTask/deploy后启动Task。这三块职责决定了下文要打穿的三条源码主线状态机、心跳、Slot 分配。04.2 源码侦探Source Sleuth04.2.1 doRunwhile(true) 自旋状态机的入口Task.run()只是包一层 MDC 上下文后调doRun()真正的状态机在doRun里// flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:581privatevoiddoRun(){// Initial State transitionwhile(true){ExecutionStatecurrentthis.executionState;if(currentExecutionState.CREATED){if(transitionState(ExecutionState.CREATED,ExecutionState.DEPLOYING)){break;// 成功进入 DEPLOYING开始干活}}elseif(currentExecutionState.FAILED){notifyFinalState();return;// 已经被外部标记失败}elseif(currentExecutionState.CANCELING){if(transitionState(ExecutionState.CANCELING,ExecutionState.CANCELED)){notifyFinalState();return;// 已经被取消}}else{thrownewIllegalStateException(Invalid state for beginning of operation of task this);}}// ... 后续加载用户代码、注册网络、恢复状态、invoke}关键源码事实doRun用一个while(true)自旋处理启动瞬间的状态竞态flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:585-616。任务线程刚启动时this.executionState可能已经被外部并发改成了FAILED或CANCELING例如调度器在任务线程启动前就 cancel 了它。自旋 transitionState的 CAS 保证要么成功抢到CREATED→DEPLOYING的迁移开始干活要么发现已被并发取消/失败而立即收尾。这段代码之所以是while(true)而不是if正是因为compareAndSet失败时状态被并发改掉要重新读一次最新状态再决定走哪条分支。04.2.2 transitionStateCAS 是状态机唯一的闸所有状态迁移都收敛到一个私有方法transitionState而它的实现核心只有一行compareAndSet// flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:1098privatebooleantransitionState(ExecutionStatecurrentState,ExecutionStatenewState,Throwablecause){if(STATE_UPDATER.compareAndSet(this,currentState,newState)){// 迁移成功记日志、记录 failureCause// ...}// ...}而STATE_UPDATER是一个AtomicReferenceFieldUpdaterflink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:155-156// flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:155privatestaticfinalAtomicReferenceFieldUpdaterTask,ExecutionStateSTATE_UPDATERAtomicReferenceFieldUpdater.newUpdater(Task.class,ExecutionState.class,executionState);关键源码事实Task的状态字段用AtomicReferenceFieldUpdater而非synchronized或AtomicReference包一层flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:155-156。这是性能与语义的双重选择状态字段是volatile的newUpdater要求字段volatilecompareAndSet提供仅当仍是旧值时迁移的原子语义保证RUNNING→FINISHED和RUNNING→FAILED这种互斥迁移不可能同时成功。synchronized会引入锁竞争AtomicReference包装则多一层对象分配AtomicReferenceFieldUpdater直接在字段上做 CAS零额外分配——这正是 Flink 在热路径上能省一次分配就省一次的体现。04.2.3 主链路的两次关键迁移DEPLOYING→INITIALIZING→RUNNING状态机进入DEPLOYING后任务要加载用户代码、注册网络、恢复状态最后才invoke。这两次迁移在restoreAndInvoke里// flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:925privatevoidrestoreAndInvoke(TaskInvokablefinalInvokable)throwsException{if(!transitionState(ExecutionState.DEPLOYING,ExecutionState.INITIALIZING)){thrownewCancelTaskException();}taskManagerActions.updateTaskExecutionState(newTaskExecutionState(executionId,ExecutionState.INITIALIZING));executingThread.setContextClassLoader(userCodeClassLoader.asClassLoader());runWithSystemExitMonitoring(finalInvokable::restore);// 恢复状态if(!transitionState(ExecutionState.INITIALIZING,ExecutionState.RUNNING)){thrownewCancelTaskException();}taskManagerActions.updateTaskExecutionState(newTaskExecutionState(executionId,ExecutionState.RUNNING));runWithSystemExitMonitoring(finalInvokable::invoke);// 执行用户代码// ...}关键源码事实restoreAndInvoke里两次迁移DEPLOYING→INITIALIZING与INITIALIZING→RUNNING每次都紧跟一次taskManagerActions.updateTaskExecutionStateflink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:929-947。INITIALIZING是恢复状态的专属状态对应ExecutionState里INITIALIZING的注释Restoring last possible valid stateRUNNING才是开始跑用户代码。关键细节finalInvokable::restore和finalInvokable::invoke都包在runWithSystemExitMonitoring里——这是 Flink 的任务线程自杀监控一旦检测到任务线程卡死如用户代码死循环会触发 JVM 退出而非无限僵持。04.2.4 正常收尾RUNNING→FINISHED 与异常分支用户代码正常跑完后任务先finish所有输出分区再尝试迁移到FINISHED// flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:782// try to mark the task as finished// if that fails, the task was canceled/failed in the meantimeif(!transitionState(ExecutionState.RUNNING,ExecutionState.FINISHED)){thrownewCancelTaskException();}而异常分支用一个while(true)自旋分类处理flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:800-835RUNNING/INITIALIZING/DEPLOYING状态下抛异常若异常是CancelTaskException则迁到CANCELED否则迁到FAILEDCANCELING状态则迁到CANCELED。这个分类逻辑确保用户代码的失败和外部取消最终落到不同的终态。04.2.5 notifyFinalState最终状态的回传收敛无论走到哪个终态最后都收敛到notifyFinalState// flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:1069privatevoidnotifyFinalState(){taskManagerActions.updateTaskExecutionState(newTaskExecutionState(executionId,this.executionState));}关键源码事实notifyFinalState把this.executionState当前值一定是三个终态之一打包成TaskExecutionState回传给 JobManagerflink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:1069-1071。doRun里FAILED分支:594、CANCELING分支:603和正常路径:869三处都调它最终汇聚到同一份回传逻辑。JobMaster 侧ExecutionGraph.updateStateflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:209收到后更新逻辑侧的Execution状态触发调度器的onTaskFinished/onTaskFailed回调形成闭环。04.2.6 TaskExecutor 的反向寻址与心跳TaskExecutor 启动时是边缘主动找中心而非等中心来连// flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:470publicvoidonStart()throwsException{startTaskExecutorServices();startRegistrationTimeout();}// flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:484privatevoidstartTaskExecutorServices()throwsException{resourceManagerLeaderRetriever.start(newResourceManagerLeaderListener());// :487taskSlotTable.start(newSlotActionsImpl(),getMainThreadExecutor());jobLeaderService.start(getAddress(),getRpcService(),haServices,newJobLeaderListenerImpl());// ...}关键源码事实TaskExecutoronStart的第一动作是resourceManagerLeaderRetriever.startflink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:487——通过 Leader 选举服务反向寻址当前活跃的 ResourceManager然后startRegistrationTimeout:481兜底一段时间内找不到 RM 就报错。这跟 JobManager 进程中心注入依赖的启动方式形成鲜明对照见 01 篇flink-runtime/src/main/java/org/apache/flink/runtime/dispatcher/Dispatcher.java:349TM 是分布式系统里的边缘节点它必须先找到谁是当前 RM才能上报自己的 Slot所以启动顺序是先寻址、再注册、再上报 Slot。心跳也分两路独立维护// flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:1063publicCompletableFutureVoidheartbeatFromJobManager(ResourceIDresourceID,AllocatedSlotReportallocatedSlotReport){returnjobManagerHeartbeatManager.requestHeartbeat(resourceID,allocatedSlotReport);}// flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:1069publicCompletableFutureVoidheartbeatFromResourceManager(ResourceIDresourceID){returnresourceManagerHeartbeatManager.requestHeartbeat(resourceID,null);}关键源码事实两路心跳各有一个独立的HeartbeatManagerflink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java:1063-1071。对 JM 的心跳携带AllocatedSlotReport告诉 JM 这台 TM 的 Slot 分配现状对 RM 的心跳则不携带 payload。两路超时独立——JM 超时会触发Task的 failover见:2700附近的超时处理RM 超时则导致 RM 认为该 TM 失联、回收其 Slot。这个双通道心跳设计让作业调度正确性依赖 JM与资源生命周期依赖 RM解耦。04.2.7 状态机时序图从 deploy 到 FINISHED下面这张时序图把 Task 从部署到完成的完整状态迁移串起来图示讲解这张时序图把 04.2 各小节的代码片段按时间顺序串成一条完整链路回答Task 从被部署到宣告完成经历了哪些状态、每步通知谁。链路起点是 JobMaster 的Execution.deploy()flink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java:557把TaskDeploymentDescriptor发给 TaskExecutorTaskExecutor 用startTaskThread()flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:567拉起 Task 线程线程进入doRun先 CAS 抢CREATED→DEPLOYING:588随后createUserCodeClassloader:637建立用户代码边界restoreAndInvoke里依次DEPLOYING→INITIALIZING:929→ 恢复状态 →INITIALIZING→RUNNING:941→ 执行用户代码其中INITIALIZING与RUNNING两处都回传updateTaskExecutionState:933/:946用户代码跑完RUNNING→FINISHED:784最后notifyFinalState:1069把终态回传给 JobMaster 的ExecutionGraph.updateStateflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/ExecutionGraph.java:209。注意图中每一条TASK → EXE的箭头都是一次 RPC 状态同步——这正是逻辑状态与物理状态双层模型靠消息对齐的体现。04.3 生产实践Production04.3.1 心跳与 Slot 相关配置配置项作用生产建议heartbeat.interval心跳间隔默认 10s大集群可放宽heartbeat.timeout心跳超时判定失联默认 50s须 interval 的 5 倍taskmanager.numberOfTaskSlots每 TM 的 Slot 数常设为 CPU 核数taskmanager.memory.process.sizeTM 进程总内存按状态大小 网络 buffer 预算心跳超时是假故障的头号来源如果heartbeat.timeout设得太小GC 停顿或网络抖动会让 JM 误判 TM 失联触发不必要的 failover。生产上timeout至少要是interval的 5 倍且要监控 GC 停顿是否逼近这个阈值。04.3.2 状态机相关的两个高频故障Task 停在 DEPLOYING/INITIALIZING 不前进DEPLOYING阶段卡在createUserCodeClassloader下载 JAR 慢或setupPartitionsAndGates网络内存不足INITIALIZING阶段卡在restore从 checkpoint 恢复大状态慢。定位看 TaskManager 日志里这两处的耗时。Invalid state for beginning of operation异常doRun的else分支flink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:613在任务线程启动时状态既不是CREATED/FAILED/CANCELING之一时抛出。这通常是同一 Task 被并发启动两次的信号根因往往在调度器侧重复部署。04.4 进阶向导Advanced04.4.1 双层状态模型的深层含义Execution逻辑与Task物理双层状态是 Flink 容错能力的关键。Execution可以有多次Task——失败后Execution从FAILED重新调度生成新的Task线程但ExecutionAttemptID会1对应Execution的 attempt 语义。这解释了为什么 JobManager 上能看到同一个 subtask 失败重试了 N 次而 TaskManager 上每次都是一条全新的Task生命周期。理解这层就能分清任务是重启了还是状态在恢复。04.4.2 与 Spark 的 Task 执行模型对比维度Flink TaskSpark Task状态机ExecutionState10 态CAS 自旋简化的 running/failed/successful线程模型每 Task 一个常驻线程流每 Task 一个线程批跑完即止状态同步updateTaskExecutionStateRPC 实时回传结果聚合到 Driver失联检测双通道心跳JM RM单一心跳Executor → Driver失败恢复checkpoint 状态重算血缘重算RDD lineage04.4.3 一个容易被误读的点很多人以为Task的executionState字段和ExecutionState枚举是JobManager 上的那套状态。实际上Task是 TaskManager 进程内的对象它的executionState是私有字段通过STATE_UPDATERflink-runtime/src/main/java/org/apache/flink/runtime/taskmanager/Task.java:155做本地 CAS而 JobManager 上Execution也维护一份ExecutionStateflink-runtime/src/main/java/org/apache/flink/runtime/executiongraph/Execution.java:253的getState。两份状态靠updateTaskExecutionState消息同步理论上存在短暂不一致窗口——这正是RECONCILING状态存在的意义failover 后 JM 用它与仍在运行的 Task 对账把两侧状态重新拉齐。04.5 小结与下一篇本篇打穿了 Flink 任务执行的一条主线Task 状态机的 CAS 自旋实现Task.doRun的while(true)transitionState的compareAndSet与TaskExecutor 的反向寻址 双通道心跳对 JM 携带 Slot 报告、对 RM 维持资源生命周期。理解了逻辑 Execution 状态与物理 Task 状态靠消息对齐的双层模型就理解了 Flink 容错与失联检测的根基。下一篇进入《状态管理StateBackend 与 KeyedState》切入点是 Task 里restoreAndInvoke的restore究竟恢复了什么StateBackend 如何创建 KeyedStateBackend、getOrCreateKeyedState如何按 keyGroup 定位状态实例以及 HashMap 与 RocksDB 两种后端的差异对应本系列 05 篇。