执行引擎: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/.../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 主线走一遍(不进代码)
以"一个普通任务跑完了"为例,端到端串一遍:
- 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:207 groupingBy(... executionId()))。
一个推论:Executor 无常驻循环,但有两个定时兜底循环。它们也只是"制造消息"的源头,不是执行循环本身:
| 定时循环 | 周期 | 干什么 |
|---|---|---|
executionDelayLoop | 1 秒 | 到点的暂停恢复 / 失败重试 / WaitFor 续跑,取出后照样走 lock → withExecution → toExecution |
executionSLAMonitorLoop | 1 秒 | 超时的 SLA 监视器,评估违约并推进执行 |
(依据:DefaultExecutor.java:257、:430 executionDelayLoop、:508 executionSLAMonitorLoop。)