数据截至 (上游 commit fcda16293bbd)
执行引擎: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/*MessageHandler | 8 个具体处理器,各接一类消息;都遵循"先锁执行,再处理" | executor/.../handler/ |
ExecutionStateStore | 执行级锁 + 持久化;lock(id, fn) 是所有推进的入口 | executor/.../ExecutionStateStore.java:lock |
ExecutorContext | 一次推进的"草稿纸":装当前 Execution + 要发的后续消息 | executor/.../ExecutorContext.java |
ExecutorService | 纯内核:process(...) 及全部 handle* 推进逻辑 | executor/.../ExecutorService.java:process |
FlowableUtils | 解析"顺序任务的下一个是谁" | core/src/main/java/io/kestra/core/runners/FlowableUtils.java:innerResolveSequentialNexts |
ExecutableUtils | 可执行任务(子流)怎么变成一次新执行 | core/src/main/java/io/kestra/core/runners/ExecutableUtils.java:subflowExecution |
2.3 两个处理器接口的分工
外壳里所有消费者都落在两个接口之一,区别只有一句话:这条消息是否会改动某次执行的状态。
ExecutorMessageHandler<T>:会。handle返回Optional<ExecutorContext>,外壳据此把新状态发回去(ExecutorMessageHandler.java:18)。MessageHandler<T>:不会。handle返回void,只做副作用(如触发别的流)(MessageHandler.java:14)。
2.4 主线走一遍(不进代码)
以"一个普通任务跑完了"为例,端到端串一遍:
- Worker 发来一条
WorkerTaskResult(任务 A 成功)。 DefaultExecutor.workerTaskResultQueue收到,交给WorkerTaskResultMessageHandler.handle。- 处理器
lock(执行ID)拿到这次执行的最新快照,包成ExecutorContext。 - 把 A 的结果并入执行(
addWorkerTaskResult),此时 A 变终态。 - 外壳继续:
ExecutionEventMessageHandler里调process(...)→handleNext解析出"该轮到 B 了",handleWorkerTasks把 B 打包成待发的WorkerTask。 - 锁内持久化新执行;
toExecution把 B 的WorkerJobEvent发给 Worker、把执行更新事件发回队列。 - 收工。等 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:268 groupingBy(... executionId()))。
一个推论:Executor 无常驻循环,但有两个定时兜底循环。它们也只是"制造消息"的源头,不是执行循环本身:
| 定时循环 | 周期 | 干什么 |
|---|---|---|
executionDelayLoop | 1 秒 | 到点的暂停恢复 / 失败重试 / WaitFor 续跑,取出后照样走 lock → withExecution → toExecution |
executionSLAMonitorLoop | 1 秒 | 超时的 SLA 监视器,评估违约并推进执行 |
(依据:DefaultExecutor.java:340、: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:180-208):
| # | 方法(ExecutorService.java) | 它负责推进什么 |
|---|---|---|
| 1 | handleRestart:922 | RESTARTED → RUNNING,并计数"重启" |
| 2 | handleEnd:894 | 顶层任务全终态时,算出 flow 的终态并渲染 outputs(转 onEnd:396) |
| 3 | handleCreatedKilling:851 | KILLING 下,把还没启动的 CREATED 子任务直接判 KILLED |
| 4 | handleKilling:940 | 所有 taskRun 都终态后,KILLING → KILLED |
| 5 | handleNext:459 | 解析下一个要跑的顺序任务(靠 FlowableUtils,见 3.4) |
| 6 | handleAfterExecution:874 | 执行进入终态后,派生 afterExecution 钩子任务(强制执行) |
| 7 | handleWorkerTasks:957 | 把 CREATED 的 runnable taskRun 打包成待发 WorkerTask;处理断点、Worker 队列路由 |
| 8 | handleFlowableTasks:497 | 处理 Flowable(Parallel/ForEach/Loop/Pause/WaitFor…)的子任务派生、重试、暂停延时 |
| 9 | handleExecutionUpdatingTasks:1229 | 就地改执行本身的任务(如 Kill、改 Labels) |
| 10 | handleExecutableTasks:1138 | 可执行任务(Subflow/ForEachItem)创建子执行(见 3.5) |
第 5 步有个短路:只有当执行不在 KILLING/KILLED/QUEUED 时才解析下一个任务——正在被杀或排队的执行不该再往前铺任务(ExecutorService.java:188)。
3.3 ExecutorContext:只加不减的草稿纸
要解决的小问题: 十道工序各自发现"要发点消息给别人"(派给 Worker、发起子流、登记延时…),怎么攒起来?
思路: ExecutorContext 是个纯累加器。它持有当前 Execution,外加几条初始容量为 0 的列表;每个 handle* 用 withXxx(...) 往里追加,方法返回 this,链式改写。process 结束后,外壳再把这些列表一次性排空成真正的队列消息。
它累加哪几类"待办":
| 列表 | 攒的是 | 谁来排空(在 ExecutionEventMessageHandler) |
|---|---|---|
nexts | 下一批要跑的 TaskRun | 合并进执行(onNexts) |
workerTasks | 要发给 Worker 的 WorkerTask | emit 到 workerJobEventQueue(:211) |
executionDelays | 定时(暂停/重试)登记 | 存进 executionDelayStateStore(:265) |
subflowExecutions | 要发起的子执行 | emit 到 executionQueue(:293) |
subflowExecutionResults | 通知父执行的子流结果 | emit 到 subflowExecutionResultQueue(:258) |
loopExecutions | Loop 的每次迭代执行 | 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:67)。
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:507)。注意规则 2/3 天然处理了并行 Flowable——只要还有子任务没终态,就不往下走,达成"等所有分支收敛"。