跳到主要内容

执行引擎:Executor 状态机

30 秒导读: Kestra 的 Executor 是把一次工作流执行(Execution)一步步往前推的大脑。它最反直觉、也最聪明的地方是:它自己不记任何状态、也没有 while 循环。执行的全部状态存在数据库里,Executor 只是一个纯函数——收到一条队列消息,就加锁载入这次执行的最新快照,跑一遍 process(...) 把它推进一格,把"接下来该发生的事"打包成新消息发出去,然后收工。推进不靠循环,靠消息不断回灌自己。

本章讲透工程含量最高的一支:Execution 如何被一步步推进。不涉及任务在 Worker 内实际运行(见 任务执行:Worker 与 RunContext)。领域模型(Flow/Task/Execution/State)见 领域模型;队列与持久化骨架见 消息骨架


1. 这是什么(零基础也能懂)

一句话定义: Executor 是工作流的状态机推进器——它接收关于某次执行的事件(任务完成了、被要求 kill 了、暂停到点了),计算这次执行的下一个状态,并派发出接下来要做的活。

它解决什么问题。 一次工作流执行是一个会持续几秒到几天的长活。中间任务在别的机器(Worker)上跑,随时可能成功、失败、被杀。谁来判断"任务 A 完了,该轮到 B 了""所有任务都终态了,整条流程算成功还是失败"?就是 Executor。

最反直觉的设计:无状态。 你可能以为 Executor 里有个大循环,盯着每个执行不放。没有。 它是一个 @Singleton 单例,不在内存里保存任何一次执行的进度。

用一个类比建立心智模型:

  • 别把 Executor 想成"管家"(记着每个客人住哪间、进度如何)。
  • 把它想成"接线员 + 一本存在数据库里的账本":每来一个电话(消息),接线员翻出对应那一页账本(加锁读 DB),照规则改一笔(推进一步),记回账本,再打几个后续电话(发消息),然后挂断。接线员脑子里什么都不留。

这个设计买到了什么:

  • 可水平扩展:因为没有内存态,可以起多个 Executor 实例分担消息,谁抢到锁谁处理。
  • 可崩溃恢复:实例挂了,状态还在 DB;换个实例接着从消息里推进。
  • 好推理:一次推进 = 一个纯函数 (Execution + 消息) → (Execution' + 待发消息),输入输出清清楚楚。

本章其余部分就是把这句"纯函数怎么推进一步"拆开讲。


2. 顶层全景(它大概怎么转)

2.1 消息驱动的外壳

Executor 分两层:外壳负责收发消息、加锁、持久化;内核 ExecutorService.process(...) 负责纯粹的状态推进。先看外壳如何把一条消息变成一次推进。

怎么读下面这张图:从上到下是一条消息的一生;关键在于——每条消息只驱动一次 process,process 之外没有循环,推进全靠消息不断回灌到顶部。

┌────────── 队列(消息骨架,见 05 章)──────────┐
各类事件 ──┤ executionEvent · workerTaskResult · killed · │
│ subflowResult · loopEvent · executionCommand │
└───────────────────────┬──────────────────────┘
│ 一条消息

DefaultExecutor.<xxx>Queue() ← 队列回调(外壳)


handler.handle(msg): lock(executionId) ← 执行级锁
│ 载入最新 Execution 快照

new ExecutorContext(execution, flow) ← 一次推进的草稿纸


ExecutorService.process(ctx) ← 推进一步(纯内核)
│ 产出:新 Execution + 一叠"待发消息"

锁内持久化 + toExecution():把待发消息发回各队列

└──► 触发下一条消息…(回到顶部)

2.2 部件一句话职责

部件干什么在哪(文件:符号)
DefaultExecutor消息外壳:订阅所有队列、跑两个定时循环、把结果发回队列executor/src/main/java/io/kestra/executor/DefaultExecutor.java:run
ExecutorMessageHandler<T>会推进执行的消息处理器接口(返回 Optional<ExecutorContext>)executor/.../ExecutorMessageHandler.java:handle
handler/*MessageHandler8 个具体处理器,各接一类消息;都遵循"先锁执行,再处理"executor/.../handler/
ExecutionStateStore执行级锁 + 持久化;lock(id, fn) 是所有推进的入口executor/.../ExecutionStateStore.java:lock
ExecutorContext一次推进的"草稿纸":装当前 Execution + 要发的后续消息executor/.../ExecutorContext.java
ExecutorService纯内核:process(...) 及全部 handle* 推进逻辑executor/.../ExecutorService.java:process
FlowableUtils解析"顺序任务的下一个是谁"core/.../runners/FlowableUtils.java:innerResolveSequentialNexts
ExecutableUtils可执行任务(子流)怎么变成一次新执行core/.../runners/ExecutableUtils.java:subflowExecution

2.3 两个处理器接口的分工

外壳里所有消费者都落在两个接口之一,区别只有一句话:这条消息是否会改动某次执行的状态

  • ExecutorMessageHandler<T>:会。handle 返回 Optional<ExecutorContext>,外壳据此把新状态发回去(ExecutorMessageHandler.java:18)。
  • MessageHandler<T>:不会。handle 返回 void,只做副作用(如触发别的流)(MessageHandler.java:14)。

2.4 主线走一遍(不进代码)

以"一个普通任务跑完了"为例,端到端串一遍:

  1. Worker 发来一条 WorkerTaskResult(任务 A 成功)。
  2. DefaultExecutor.workerTaskResultQueue 收到,交给 WorkerTaskResultMessageHandler.handle
  3. 处理器 lock(执行ID) 拿到这次执行的最新快照,包成 ExecutorContext
  4. 把 A 的结果并入执行(addWorkerTaskResult),此时 A 变终态。
  5. 外壳继续:ExecutionEventMessageHandler 里调 process(...)handleNext 解析出"该轮到 B 了",handleWorkerTasks 把 B 打包成待发的 WorkerTask
  6. 锁内持久化新执行;toExecution 把 B 的 WorkerJobEvent 发给 Worker、把执行更新事件发回队列。
  7. 收工。等 B 完成时,又一条 WorkerTaskResult 回灌,重复。

3. 核心原理(逐个机制,由浅入深)

3.1 无状态 + 一把锁 = 一次推进

要解决的小问题: 多个 Executor 实例、多条消息可能同时碰到同一次执行,怎么不打架?

思路: 不在内存里存执行态,而是每次推进都从"执行级锁"里取最新快照、改完立刻在锁内写回。锁的粒度是单个 executionId,所以不同执行天然并行,同一执行天然串行。

ExecutionStateStore.lock 的签名就把这个模式钉死了——它要你传一个"从旧执行算出新 ExecutorContext"的纯函数:

// executor/src/main/java/io/kestra/executor/ExecutionStateStore.java:16
Optional<ExecutorContext> lock(String executionId, Function<Execution, ExecutorContext> function);

DefaultExecutor 甚至在批处理层面再加一道保险:同一批 executionEvent 先按 executionId 分组,同组顺序处理、不同组并发,避免同一执行的多条消息挤在一起(DefaultExecutor.java:207 groupingBy(... executionId()))。

一个推论:Executor 无常驻循环,但有两个定时兜底循环。它们也只是"制造消息"的源头,不是执行循环本身:

定时循环周期干什么
executionDelayLoop1 秒到点的暂停恢复 / 失败重试 / WaitFor 续跑,取出后照样走 lock → withExecution → toExecution
executionSLAMonitorLoop1 秒超时的 SLA 监视器,评估违约并推进执行

(依据:DefaultExecutor.java:257:430 executionDelayLoop:508 executionSLAMonitorLoop。)

3.2 process 流水线:一次推进的十道工序

要解决的小问题: "推进一步"到底做哪些判断?

思路: process 是一条固定顺序的流水线,依次调十来个 handle*,每个只管一种推进,顺序本身编码了状态机的优先级(先处理重启/收尾/杀,再处理正常派发)。

先看它的门禁与骨架:

// executor/src/main/java/io/kestra/executor/ExecutorService.java:135
public ExecutorContext process(ExecutorContext executor) {
// 前置失败 / 并发限流已终结 / 已终态 → 原样返回,不推进
if (!executor.canBeProcessed() || executionService.isTerminated(executor.getFlow(), executor.getExecution())) {
return executor;
}
try {
executor = this.handleRestart(executor);
executor = this.handleEnd(executor);
// ...(见下表)
} catch (Exception e) {
return executor.withException(e, "process"); // 不崩:标记异常,留待收尾失败
}
return executor;
}

十道工序按调用顺序(ExecutorService.java:144-172):

#方法(ExecutorService.java)它负责推进什么
1handleRestart:922RESTARTEDRUNNING,并计数"重启"
2handleEnd:894顶层任务全终态时,算出 flow 的终态并渲染 outputs(转 onEnd:396)
3handleCreatedKilling:851KILLING 下,把还没启动的 CREATED 子任务直接判 KILLED
4handleKilling:940所有 taskRun 都终态后,KILLINGKILLED
5handleNext:459解析下一个要跑的顺序任务(靠 FlowableUtils,见 3.4)
6handleAfterExecution:874执行进入终态后,派生 afterExecution 钩子任务(强制执行)
7handleWorkerTasks:957CREATED 的 runnable taskRun 打包成待发 WorkerTask;处理断点、Worker 队列路由
8handleFlowableTasks:497处理 Flowable(Parallel/ForEach/Loop/Pause/WaitFor…)的子任务派生、重试、暂停延时
9handleExecutionUpdatingTasks:1229就地改执行本身的任务(如 Kill、改 Labels)
10handleExecutableTasks:1138可执行任务(Subflow/ForEachItem)创建子执行(见 3.5)

第 5 步有个短路:只有当执行不在 KILLING/KILLED/QUEUED 时才解析下一个任务——正在被杀或排队的执行不该再往前铺任务(ExecutorService.java:152)。

3.3 ExecutorContext:只加不减的草稿纸

要解决的小问题: 十道工序各自发现"要发点消息给别人"(派给 Worker、发起子流、登记延时…),怎么攒起来?

思路: ExecutorContext 是个纯累加器。它持有当前 Execution,外加几条初始容量为 0 的列表;每个 handle*withXxx(...) 往里追加,方法返回 this,链式改写。process 结束后,外壳再把这些列表一次性排空成真正的队列消息。

它累加哪几类"待办":

列表攒的是谁来排空(在 ExecutionEventMessageHandler)
nexts下一批要跑的 TaskRun合并进执行(onNexts)
workerTasks要发给 Worker 的 WorkerTaskemit 到 workerJobEventQueue(:211)
executionDelays定时(暂停/重试)登记存进 executionDelayStateStore(:265)
subflowExecutions要发起的子执行emit 到 executionQueue(:293)
subflowExecutionResults通知父执行的子流结果emit 到 subflowExecutionResultQueue(:258)
loopExecutionsLoop 的每次迭代执行emit 到 executionQueue(:299)

withExecution 顺手做了一件精巧的事:记录本轮走过的每一个不同状态。因为一次推进可能一步跨好几个状态(如 PAUSED → RUNNING → SUCCESS),它把每次真正的状态变化 append 进 stateTransitions,好让收尾时每个中间态都能触达 flow 触发器:

// executor/src/main/java/io/kestra/executor/ExecutorContext.java:72
public ExecutorContext withExecution(Execution execution, String from) {
this.execution = execution;
this.from.add(from);
this.executionUpdated = true;
State.Type newState = execution.getState().getCurrent();
if (!newState.equals(stateTransitions.getLast())) { // 只在状态真变了时记一笔
stateTransitions.add(newState);
}
return this;
}

from 这个字符串列表纯为调试:每次改写都留个来源标签("handleRestart""onNexts"…),日志里能一眼看出这轮推进被哪些工序动过。

canBeProcessed() 则是 3.2 那道门禁的实现——执行已删除 / 暂停 / 断点 / 排队 / flow 无效,都不推进(ExecutorContext.java:61)。

3.4 Flowable 的 next 解析:顺序任务谁是下一个

要解决的小问题: 给定一串顺序任务和"已经跑到哪了",算出下一个该创建的是谁——还是"都跑完了,别再派了"。

思路: 这是纯计算,抽在 FlowableUtils.innerResolveSequentialNexts。规则只有三条,按序判断:

1) 一个 taskRun 都还没建 ──► 派第一个任务
2) 有任一 CREATED/SUBMITTED/RUNNING ──► 什么都别派(还在跑,等着)
3) 否则找"最后一个终态"的任务,派它的下一个

真源码:

// core/src/main/java/io/kestra/core/runners/FlowableUtils.java:77
private static List<NextTaskRun> innerResolveSequentialNexts(Execution execution, List<ResolvedTask> currentTasks, TaskRun parentTaskRun) {
if (currentTasks == null || currentTasks.isEmpty() || execution.getState().getCurrent() == State.Type.KILLING) {
return Collections.emptyList();
}
List<TaskRun> taskRuns = execution.findTaskRunByTasks(currentTasks, parentTaskRun);
if (taskRuns.isEmpty()) { // 规则 1
return Collections.singletonList(currentTasks.getFirst().toNextTaskRun(execution));
}
if (taskRuns.stream().anyMatch(t -> t.getState().isCreated() // 规则 2
|| t.getState().getCurrent() == State.Type.SUBMITTED || t.getState().isRunning())) {
return Collections.emptyList();
}
Optional<TaskRun> lastTerminated = execution.findLastTerminated(taskRuns); // 规则 3
// ... 找到 lastTerminated 在列表里的下标,返回下标+1 的任务
}

handleNext 就是调它:普通执行解析整条 flow 的顺序任务,LOOP 类型的执行则只解析这个 Loop 自己的子任务(ExecutorService.java:459)。注意规则 2/3 天然处理了并行 Flowable——只要还有子任务没终态,就不往下走,达成"等所有分支收敛"。

3.5 可执行子流派发:一个任务变成一次新执行

要解决的小问题: Subflow/ForEachItem 这类"可执行任务"不在 Worker 上跑代码,而是要发起一次全新的执行。怎么把父执行里的一个 taskRun 变成一条子执行?

思路: 这类任务实现 ExecutableTaskhandleWorkerTasks 会先把它当普通 worker task 打包进 workerTasks;handleExecutableTasks 再从 workerTasks挑出并移除它们,创建子执行、并登记"回头通知父执行"的结果。

handleExecutableTasksremoveIf 遍历,命中 ExecutableTask 就处理并从待发列表移走(ExecutorService.java:1142):

  • 先把该 taskRun 标 RUNNING,避免失败时重复尝试(:1155)。
  • runIf,不成立就标 SKIPPED 收工(:1163)。
  • executableTask.createSubflowExecutions(...) 造子执行;若为空,直接 SUCCESS(:1176)。
  • 否则把子执行塞进 executions,并按"是否等待子执行"决定父 taskRun 现在就 SUCCESS 还是挂起等结果(:1193)。

真正把"一次子执行长什么样"拼出来的是 ExecutableUtils.subflowExecution:它渲染子流的 namespace/id/revision、算好继承的 labels、把父执行的 executionId/taskRunId/... 塞进子执行的 ExecutionTrigger.variables 里当"回执地址"(ExecutableUtils.java:209)。子执行终态时,外壳靠 ExecutableUtils.isSubflow 认出它、按这些变量找回父执行发 SubflowExecutionEnd(ExecutableUtils.java:313;发送处 DefaultExecutor.java:635)。

一个易错点已在代码里点明:子执行永远不是 LOOP kind——LOOP 只留给 Loop 的虚拟迭代执行,子流是独立执行(ExecutableUtils.java:232)。

3.6 容错:Flowable 不许崩,失败要收得住

要解决的小问题: 推进过程中出异常(渲染失败、解析状态失败、找不到 flow),不能让整个 Executor 挂掉,也不能让执行永远卡住。

思路: 多层"就地降级",核心一句是——宁可把这次执行判 FAILED,也不让线程崩

三处典型防线:

  • process 顶层 try/catch:任一 handle* 抛异常,process 不再往下,而是 withException(e, "process") 记下异常;外壳 toExecution 看到异常就走 handleFailedExecutionFromExecutor 把执行判失败(ExecutorService.java:173:228;调用点 DefaultExecutor.java:586)。
  • Flowable 状态解析的 panic 模式:childWorkerTaskResult 里若 resolveState 抛错,注释直言"Flowable 不该失败,这是一种 panic 模式",于是兜底 state = FAILED 让下一步能继续,而不是卡住(ExecutorService.java:244)。
  • 单个 Flowable 任务的 try/catch:handleFlowableTasks 逐任务包裹,某个 Flowable 处理炸了就把它标 FAILED,并把它待发的 WorkerTask 一并替换成 FAILED,防止它又转 RUNNING(ExecutorService.java:704)。

还有一处专门防死循环:当异常是 FlowNotFoundException(找不到 flow),处理器会先看执行是不是已经 FAILED——是就直接放弃,不再反复判失败,免得同一条消息把执行永远反复处理(WorkerTaskResultMessageHandler.java:71ExecutionEventMessageHandler.java:313)。

最后,onNexts 是"派新任务"时的小状态机:把新 taskRun 并入执行;若执行还在 CREATED,顺手翻成 RUNNING 并记一条"Flow started"(ExecutorService.java:211)。


4. 深入实现(几个 handle* 走读)

这一节给要读源码的人,补几个上面点到、但值得看细节的工序。

4.1 handleWorkerTasks:打包、路由、断点

它做三件事(ExecutorService.java:957):

  1. 打包:把每个 CREATED 且没有测试夹具的 taskRun,连同其 RunContext 造成 WorkerTask;注入 OpenTelemetry 的 traceparent(:977)。
  2. Worker 队列路由:调 workerQueueService.resolveWorkerQueueForJob。找不到合适队列时按 disposition 处理——FAIL 直接失败该 taskRun、CANCELCANCELLEDWAIT_AND_DISPATCH 记日志后照发(:986:1000)。
  3. 断点:若某 CREATED taskRun 命中断点,整条执行转 BREAKPOINT 并停在这(:1082shouldSuspend:1133)。

注意它不直接发消息——只把结果攒进 ExecutorContext.workerTasks。真正 emit 到 Worker 队列、并把 taskRun 从 CREATED 推到 SUBMITTED 的,是外壳 ExecutionEventMessageHandler(:213)。这正是"内核纯计算、外壳做 IO"的分工。

4.2 外壳如何排空 ExecutorContext

ExecutionEventMessageHandler.handle 是把内核产物变成真实消息的地方,顺序也有讲究(ExecutionEventMessageHandler.java:88):

  1. 处理排程延迟、SLA 监视器、并发限流(在 process 之前)(:109:165)。
  2. executorService.process(executor)(:174)。
  3. onNexts 合并新任务(:176)。
  4. 遍历 workerTasks:求 runIf → 发 WorkerJobEvent 给 Worker、taskRun 转 SUBMITTED;Flowable 则就地转 RUNNING(:184:251)。
  5. 依次排空 subflowExecutionResultsexecutionDelayssubflowExecutionsloopExecutions(:254:300)。

4.3 toExecution:收尾与"每个状态都触发一次"

DefaultExecutor.toExecution 是每次推进的出口(DefaultExecutor.java:582)。除了把执行更新/终态事件发回队列,它做了一件容易忽略但重要的事——为本轮走过的每一个中间状态各触发一次 flow 触发器:

// executor/src/main/java/io/kestra/executor/DefaultExecutor.java:626
List<State.Type> transitions = executor.getStateTransitions();
for (int i = 1; i < transitions.size(); i++) { // 跳过 index 0(入口态)
State.Type transitionState = transitions.get(i);
processFlowTriggers(transitionState == execution.getState().getCurrent() ? execution : execution.withState(transitions.get(i)));
}

这就是 3.3 里 stateTransitions 的用途:一次推进哪怕从 PAUSED 一路冲到 SUCCESS,每个中间态都不会漏掉对下游"按状态监听"的流的通知。终态时它还负责发 SubflowExecutionEnd、发 Loop 迭代事件、按并发策略弹出排队执行(:633:737)。

State.Type 的全部取值见 core/.../models/flows/State.java:239(CREATED/SUBMITTED/RUNNING/PAUSED/RESTARTED/KILLING/SUCCESS/WARNING/FAILED/KILLED/CANCELLED/QUEUED/RETRYING/RETRIED/SKIPPED/BREAKPOINT/RESUBMITTED)。


5. 巧妙之处(可借鉴的技术)

  • 无状态 + 执行级锁 = 天然可扩展的状态机。 把状态外置到 DB、把推进写成纯函数 lock(id, Execution → ExecutorContext),一举拿下水平扩展、崩溃恢复、易推理三件事(ExecutionStateStore.java:16)。
  • 累加器模式隔离计算与 IO。 内核 process 只往 ExecutorContext 里追加"待办",绝不碰队列;外壳统一排空。测试内核时不需要真队列(ExecutorContext.java;排空处 ExecutionEventMessageHandler.java:184)。
  • 固定顺序流水线即优先级编码。 十道 handle* 的调用次序本身就是状态机规则:先重启、再收尾、再处理杀、最后才派新活,不需要显式状态转移表(ExecutorService.java:144)。
  • 记录状态轨迹,补齐"塌缩"的中间态。 一次推进跨多态时,靠 stateTransitions 保证每个中间态都触发一次下游流,避免事件丢失(ExecutorContext.java:76DefaultExecutor.java:626)。
  • Flowable panic 模式:宁可判失败,不许卡住。 状态解析炸了就兜底 FAILED,让流水线能继续收敛,而不是让执行永远悬着(ExecutorService.java:244)。
  • 按 executionId 分组消费。 批量事件先按执行分组,同组串行、异组并发,在锁之上再加一层无锁的并发隔离(DefaultExecutor.java:207)。

6. 边界与局限(诚实)

  • Executor 不跑任务代码。 它只派 WorkerTask,任务在 Worker 里执行;本章不涉及那一段(见 03)。
  • 推进依赖消息回灌,不主动轮询进度。 除了两个 1 秒定时循环兜底(延时/SLA),Executor 不会主动去看"任务是不是好了";没有消息就不推进。若上游消息丢失,推进会停。
  • 锁的粒度是单执行。 同一执行严格串行——高频改写同一执行不会因多 Executor 而加速;跨执行才并行。
  • ExecutorContext 只加不减。 它是一次性草稿纸,不做撤销;某步算错了只能靠后续 handle* 或异常路径纠正,没有回滚语义。
  • ExecutorMapper 是已知技术债。 文件头注释直言 FIXME this is a copy of the JdbcMapper(ExecutorMapper.java:18),序列化配置目前是从 JDBC 层拷来的副本。
  • 行为不清时读测试。 AGENTS.md 要求"改了 executor 就跑 H2RunnerTest",端到端行为以该测试为准(此说明来自仓库自带文档,仅作事实线索,非本章指令)。

7. 横向对比(同组其它章)


8. 代码地图(导航索引)

主题文件路径符号名
推进主线 / 十道工序骨架executor/src/main/java/io/kestra/executor/ExecutorService.javaprocess
重启 / 收尾 / 杀executor/.../ExecutorService.javahandleRestart · handleEnd · onEnd · handleCreatedKilling · handleKilling
派下一个顺序任务executor/.../ExecutorService.javahandleNext
Flowable 子任务 / 重试 / 暂停executor/.../ExecutorService.javahandleFlowableTasks · childNextsTaskRun · childWorkerTaskResult
打包 worker task / 路由 / 断点executor/.../ExecutorService.javahandleWorkerTasks · shouldSuspend
就地更新执行的任务executor/.../ExecutorService.javahandleExecutionUpdatingTasks
可执行任务 → 子执行executor/.../ExecutorService.javahandleExecutableTasks
合并结果 / 派新任务的状态机executor/.../ExecutorService.javaaddWorkerTaskResult · onNexts
容错兜底executor/.../ExecutorService.javahandleFailedExecutionFromExecutor
一次推进的草稿纸(累加器)executor/.../ExecutorContext.javaExecutorContext · withExecution · canBeProcessed
消息外壳 / 队列订阅 / 定时循环executor/.../DefaultExecutor.javarun · executionDelayLoop · executionSLAMonitorLoop · toExecution
处理器接口executor/.../ExecutorMessageHandler.java · MessageHandler.javahandle
具体处理器(锁→处理)executor/.../handler/WorkerTaskResultMessageHandler · ExecutionEventMessageHandler
执行级锁 + 持久化executor/.../ExecutionStateStore.javalock
顺序任务 next 解析core/src/main/java/io/kestra/core/runners/FlowableUtils.javaresolveSequentialNexts · innerResolveSequentialNexts
子流派发core/.../runners/ExecutableUtils.javasubflowExecution · isSubflow
状态枚举core/.../models/flows/State.javaState.Type