数据截至 (上游 commit fcda16293bbd)
Kestra — 架构与原理
30 秒导读: Kestra 是一个开源的工作流编排平台——你用几行 YAML 描述"先干什么、再干什么、出错怎么办",它就负责按顺序、可靠地把这些步骤跑完,失败能重试、能定时、能被事件触发。它的内核是一台队列驱动的状态机:所有组件之间不直接调用,而是往队列里丢消息、从队列里取消息,因此同一套代码既能在你笔记本上单机跑,也能在集群里扛住几百万次执行。
1. 这是什么(零基础也能懂)
一句话定义: Kestra 是一个"把一堆步骤按你写的顺序、可靠地跑完"的平台——业界叫工作流编排(workflow orchestration)。
解决什么问题 / 给谁用。 假设你是数据工程师,每天要做这么一串活:
- 从数据库导出昨天的订单;
- 用 Python 清洗一遍;
- 灌进数仓;
- 成功了发 Slack,失败了报警。
手写脚本 + crontab 能跑,但一旦某步失败、要重试、要看历史、要在出错时通知人,脚本就变成一团乱麻。Kestra 把这类"多步骤、有依赖、要调度、要容错"的活儿,变成一份声明式的 YAML,并配一个能看到每步状态的 UI。
它能做什么(功能):
- 声明式编排——用 YAML 描述任务、依赖、错误分支,不用写调度逻辑。
- 两种启动方式——定时(Cron)触发,或事件(文件到达、消息、Webhook)触发。
- 任意语言任务——通过插件跑 Python / Shell / SQL / Docker / 云服务等。
- 容错——重试、超时、错误处理、并发限制、暂停/恢复。
- 可视化 + 版本化——UI 里画拓扑图,同时 YAML 始终是唯一事实来源。
用起来什么样。 一份最小的 Flow 就是一段 YAML:
id: hello_world
namespace: dev
tasks:
- id: say_hello
type: io.kestra.plugin.core.log.Log
message: "Hello, World!"
id 是流程名,tasks 是要依次执行的步骤,每个 type 指向一个插件类(全限定 Java 类名)。把它贴进 UI 点运行,就能看到一次执行(Execution)从 CREATED 一路走到 SUCCESS(README.md:150-158)。
一句话直觉/类比。 把 Kestra 想成一个流水线上的调度员:你写的 YAML 是"工艺卡片",调度员照着卡片一步步发号施令——但它自己不干活,只负责"下一步该谁上、上一步成没成",真正的体力活交给下游工人(Worker)。而且这个调度员没有记忆:每次要做决定,它都重新把整张卡片和当前进度读一遍,再决定下一步——这正是它能水平扩展、扛住海量执行的原因。
2. 顶层全景(它大概怎么转)
先看怎么读这张图:从左到右是一次执行的生命周期,中间那条队列(Queue)总线是关键——所有组件都不互相直接调用,只往队列丢消息、从队列取消息。
┌──────────┐ ①下 Create 指令 ┌──────────────────────────────┐
│ Web 层 │ ───────────────────────▶│ │
│(UI/API) │ │ Queue 消息总线 │
└──────────┘ │ (executionCommand / │
┌──────────┐ ①凭触发器造执行 │ execution / workerTaskResult │
│ Scheduler │ ───────────────────────▶│ / … 各种队列) │
│(定时/事件)│ │ │
└──────────┘ └──────────────────────────────┘
▲ ③回传结果 │ ②取执行
│ ▼
┌──────────┐ ┌──────────────┐
│ Worker │◀──────│ Executor │
│(跑任务) │ ②派活 │(无状态状态机) │
└──────────┘ └──────────────┘
│
▼
仓储/状态存储
(H2 / Postgres / MySQL)
主线走一遍(高层,不进代码):
- 入口——不管是你在 UI 点"运行",还是 Scheduler 定时触发,最终都往
executionCommandQueue丢一条 ExecutionCommand(例如Create),而不是直接造执行。Web 层见ExecutionController;Scheduler 侧见DefaultTriggerExecutionPublisher,两条入口汇到同一个命令队列(webserver/…/ExecutionController.java:632-640;core/…/DefaultTriggerExecutionPublisher.java:26)。 - 推进——Executor 从队列取出执行,跑一遍状态机
process(...):算出"下一个该跑的任务",把它作为 WorkerTask 发给 Worker;然后把更新后的执行写回队列/仓储,自己不留任何内存状态(executor/…/ExecutorService.java:135;executor/…/DefaultExecutor.java:200)。 - 执行——Worker 收到 WorkerTask,真正调用那个插件的
run(RunContext),跑完把 WorkerTaskResult 发回队列(worker/…/AbstractWorker.java:120;core/…/RunnableTask.java)。 - 收敛——Executor 收到结果,再跑一次状态机,推进到下一个任务,直到没有下一步 → 执行进入终态(
SUCCESS/FAILED/ …)。
部件一句话职责:
| 部件 | 干什么 | 在哪个模块 / 类 |
|---|---|---|
| Web 层 | 把 HTTP 请求(建流程、跑执行、Webhook)变成命令消息 | webserver · ExecutionController / FlowController |
| Scheduler | 按 Cron / 事件触发器凭空造出执行命令 | scheduler · DefaultScheduler / TriggerEventHandler |
| Executor | 无状态状态机,决定"下一步跑什么",推进执行 | executor · ExecutorService / DefaultExecutor |
| Worker | 真正执行单个任务,回传结果 | worker · AbstractWorker / WorkerLoop |
| Queue | 组件间唯一通信方式,可插拔 | queue / queue-jdbc · QueueInterface |
| 仓储/状态存储 | 存执行、流程、触发器状态,可插拔 | jdbc-h2 / jdbc-postgres / jdbc-mysql |
| Plugin 系统 | 加载"任务类型"(type: 指向的类) | core · PluginRegistry / PluginScanner |
一个关键设计要先说破: Executor 是无状态的。它不在内存里维护"这次执行到哪了",而是每处理一条消息,就把执行加锁读出→跑一遍状态机→写回。所以状态机可以随便加机器并行,单点挂了也不丢进度——进度全在队列和仓储里。这也是"单机 vs 集群"能靠换后端实现的根因:换的只是 Queue / 仓储的实现,状态机代码一行不动(详见 消息骨架:Queue 抽象与持久化)。
3. 阅读地图(该按什么顺序读)
建议从数据结构入手,再看引擎,最后看边缘接入。六章由浅入深:
-
领域模型:Flow / Task / Execution / State —— 先搞懂四个核心结构:
Flow(你写的 YAML,含tasks/errors/triggers)、Task(一个步骤)、Execution(一次真实运行,含taskRunList)、State(状态机的状态枚举)。不懂这四个,后面全是空中楼阁。依据:core/…/models/flows/Flow.java、core/…/models/executions/Execution.java、core/…/models/flows/State.java。 -
执行引擎:Executor 状态机 —— 全项目最核心的一章。看
ExecutorService.process(...)那串handleRestart → handleEnd → handleNext → handleWorkerTasks → handleFlowableTasks → …的 handler 链,理解"无状态状态机如何靠反复重算推进执行"。依据:executor/…/ExecutorService.java:135-185。 -
任务执行:Worker 与 RunContext —— Worker 怎么把一条 WorkerTask 变成真正的进程/容器,
RunContext又如何给任务提供渲染变量、日志、存储、指标这套"运行时环境"。依据:worker/…/AbstractWorker.java、core/…/runners/RunContext.java。 -
调度与触发:Scheduler 与 Trigger —— Executor 只会"推进已有执行";执行本身从哪来?这一章讲 Scheduler 如何评估 Cron / 事件触发器,凭空生成新执行。依据:scheduler/…/DefaultScheduler.java、scheduler/…/internals/SchedulableEvaluator.java、scheduler/…/TriggerEventHandler.java:343-344。
-
消息骨架:Queue 抽象与持久化 —— 把前四章串起来的"血管"。讲
QueueInterface的 emit/receive 语义、QueueFactoryInterface列出的所有队列种类,以及 JDBC 实现如何把队列做成一张表。这里能看清"单机 H2 vs 集群 Postgres/MySQL"到底换了什么。依据:core/…/queues/QueueInterface.java、queue/…/QueueFactoryInterface.java、queue-jdbc/…/JdbcQueueFactory.java。 -
扩展与接入:插件系统与 Web 层 —— 收尾:
type:里那串类名怎么被PluginRegistry找到并实例化;Web 层怎么把 REST/Webhook 请求变成执行命令。依据:core/…/plugins/PluginRegistry.java、core/…/plugins/PluginScanner.java、webserver/…/ExecutionController.java。
4. 巧妙之处(读者要带走的精华)
这几条是 Kestra 架构里不显然、但很值得借鉴的设计决策:
-
无状态状态机 + "重算而非记忆"。
ExecutorService.process(...)每次都把整个执行重跑一遍 handler 链,而不是维护"上次算到哪"的增量状态(executor/…/ExecutorService.java:135)。代价是每步都重算,收益是任意水平扩展 + 崩溃不丢进度——因为没有任何只存在内存里的进度。 -
命令(Command)与执行(Execution)分家。 Web 层和 Scheduler 都不直接造
Execution,而是发一条ExecutionCommand(如Create)到executionCommandQueue,由 Executor 统一物化成执行(webserver/…/ExecutionController.java:632-640)。好处:所有"想让某执行发生变化"的意图(创建、重启、改标签、改状态)都走同一条命令通道,入口再多,状态变更只有一处。 -
State 是不可变、带完整历史的值对象。 每次状态变更都
new一个新State,把旧历史拷进去再追加一条(core/…/models/flows/State.java:47-51)。这让"这次执行经历过哪些状态、各在什么时刻"天然可审计,也让getDuration()之类的推导零成本。 -
同一套接口,单机/集群靠"换后端"切换。
QueueInterface与仓储接口是抽象的;server local只是把kestra.queue.type/kestra.repository.type设成h2,用内嵌数据库把"队列"做成一张表,单 JVM 跑全部组件(cli/…/LocalCommand.java:33-40)。集群则换 Postgres/MySQL、拆分服务进程。业务代码对此无感知。 -
队列种类高度细分,而非一个大 topic。
QueueFactoryInterface列出了 execution、executionCommand、workerTaskResult、kill、trigger 等十几种独立队列(queue/…/QueueFactoryInterface.java:20-60)。每类消息各走各的通道,便于分别限流、监控队列积压(queueLagForConsumerGroup)、独立扩缩。 -
插件即
type:字符串 → Java 类的映射。 YAML 里type: io.kestra.plugin.core.log.Log会被PluginRegistry.findClassByIdentifier(...)解析成真实类并实例化(core/…/plugins/PluginRegistry.java:66)。加一种新任务 = 加一个实现RunnableTask的类,内核零改动。
5. 代码地图(导航索引)
下面这张表让你/agent 直接跳进源码;用符号名 grep 比行号更抗上游漂移。路径相对克隆根,行号 as-of e4baca6。
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 流程定义(YAML 落成对象) | core/src/main/java/io/kestra/core/models/flows/Flow.java | Flow(tasks / errors / triggers 字段) |
| 一次运行 + 任务运行列表 | core/src/main/java/io/kestra/core/models/executions/Execution.java | Execution(id:71 / taskRunList:84 / state:105) |
| 单个任务的一次运行 | core/src/main/java/io/kestra/core/models/executions/TaskRun.java | TaskRun(state:63 / attempts:56) |
| 状态枚举 + 不可变状态机值 | core/src/main/java/io/kestra/core/models/flows/State.java | State.Type(:239)、State.withState(:73) |
| 执行引擎主循环(状态机) | executor/src/main/java/io/kestra/executor/ExecutorService.java | process(:135)、handleNext(:459)、handleWorkerTasks(:957) |
| 执行引擎服务(订阅队列、写回) | executor/src/main/java/io/kestra/executor/DefaultExecutor.java | run(:200)、toExecution(:578) |
| 可执行任务的契约 | core/src/main/java/io/kestra/core/models/tasks/RunnableTask.java | RunnableTask.run(RunContext) |
| Worker 主体与启动 | worker/src/main/java/io/kestra/worker/AbstractWorker.java | AbstractWorker.start(:120) |
| Worker 轮询循环骨架 | worker/src/main/java/io/kestra/worker/WorkerLoop.java | WorkerLoop.runLoop(:88) |
| 任务运行时环境 | core/src/main/java/io/kestra/core/runners/RunContext.java | render / logger / storage / metric |
| 调度器主体 | scheduler/src/main/java/io/kestra/scheduler/DefaultScheduler.java | DefaultScheduler(:45) |
| 触发器评估 | scheduler/src/main/java/io/kestra/scheduler/internals/SchedulableEvaluator.java | SchedulableEvaluator.evaluate(:36) |
| 触发器→执行 的落地 | scheduler/src/main/java/io/kestra/scheduler/TriggerEventHandler.java | toExecution + triggerExecutionPublisher.send(:343-344) |
| 触发器发布到命令队列 | core/src/main/java/io/kestra/core/scheduler/service/DefaultTriggerExecutionPublisher.java | DefaultTriggerExecutionPublisher(executionCommandQueue:26) |
| 队列抽象(emit/receive) | core/src/main/java/io/kestra/core/queues/QueueInterface.java | QueueInterface(emit:16 / receive:74) |
| 所有队列种类清单 | queue/src/main/java/io/kestra/queue/QueueFactoryInterface.java | QueueFactoryInterface(:20) |
| JDBC 队列工厂 | queue-jdbc/src/main/java/io/kestra/queue/jdbc/JdbcQueueFactory.java | JdbcQueueFactory |
| Web 层:执行命令入口 | webserver/src/main/java/io/kestra/webserver/controllers/api/ExecutionController.java | executionCommandQueue.emit(:640)、webhook(:524) |
| 插件注册表(type→类) | core/src/main/java/io/kestra/core/plugins/PluginRegistry.java | findClassByIdentifier(:66) |
| 插件扫描/加载 | core/src/main/java/io/kestra/core/plugins/PluginScanner.java | PluginScanner.scan(:59) |
| 单机启动(H2 后端) | cli/src/main/java/io/kestra/cli/commands/servers/LocalCommand.java | LocalCommand(h2 配置:33-40) |
| 单 JVM 跑全部服务 | cli/src/main/java/io/kestra/cli/commands/servers/StandAloneCommand.java | StandAloneCommand(:32) |
一句话回顾: YAML →
Flow;点运行/触发 →ExecutionCommand进队列 →Executor无状态状态机推进 →Worker跑任务回传 → 直到终态。所有连接线都是队列,而队列和仓储可换后端——这就是 Kestra 既能单机又能扛百万级的全部秘密。