跳到主要内容

Kestra — 架构与原理

30 秒导读: Kestra 是一个开源的工作流编排平台——你用几行 YAML 描述"先干什么、再干什么、出错怎么办",它就负责按顺序、可靠地把这些步骤跑完,失败能重试、能定时、能被事件触发。它的内核是一台队列驱动的状态机:所有组件之间不直接调用,而是往队列里丢消息、从队列里取消息,因此同一套代码既能在你笔记本上单机跑,也能在集群里扛住几百万次执行。


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

一句话定义: Kestra 是一个"把一堆步骤按你写的顺序、可靠地跑完"的平台——业界叫工作流编排(workflow orchestration)

解决什么问题 / 给谁用。 假设你是数据工程师,每天要做这么一串活:

  1. 从数据库导出昨天的订单;
  2. 用 Python 清洗一遍;
  3. 灌进数仓;
  4. 成功了发 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)

主线走一遍(高层,不进代码):

  1. 入口——不管是你在 UI 点"运行",还是 Scheduler 定时触发,最终都往 executionCommandQueue 丢一条 ExecutionCommand(例如 Create),而不是直接造执行。Web 层见 ExecutionController;Scheduler 侧见 DefaultTriggerExecutionPublisher,两条入口汇到同一个命令队列(webserver/…/ExecutionController.java:632-640;core/…/DefaultTriggerExecutionPublisher.java:26)。
  2. 推进——Executor 从队列取出执行,跑一遍状态机 process(...):算出"下一个该跑的任务",把它作为 WorkerTask 发给 Worker;然后把更新后的执行写回队列/仓储,自己不留任何内存状态(executor/…/ExecutorService.java:135;executor/…/DefaultExecutor.java:200)。
  3. 执行——Worker 收到 WorkerTask,真正调用那个插件的 run(RunContext),跑完把 WorkerTaskResult 发回队列(worker/…/AbstractWorker.java:120;core/…/RunnableTask.java)。
  4. 收敛——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. 阅读地图(该按什么顺序读)

建议从数据结构入手,再看引擎,最后看边缘接入。六章由浅入深:

  1. 领域模型: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。

  2. 执行引擎:Executor 状态机 —— 全项目最核心的一章。看 ExecutorService.process(...) 那串 handleRestart → handleEnd → handleNext → handleWorkerTasks → handleFlowableTasks → … 的 handler 链,理解"无状态状态机如何靠反复重算推进执行"。依据:executor/…/ExecutorService.java:135-185。

  3. 任务执行:Worker 与 RunContext —— Worker 怎么把一条 WorkerTask 变成真正的进程/容器,RunContext 又如何给任务提供渲染变量、日志、存储、指标这套"运行时环境"。依据:worker/…/AbstractWorker.java、core/…/runners/RunContext.java。

  4. 调度与触发:Scheduler 与 Trigger —— Executor 只会"推进已有执行";执行本身从哪来?这一章讲 Scheduler 如何评估 Cron / 事件触发器,凭空生成新执行。依据:scheduler/…/DefaultScheduler.java、scheduler/…/internals/SchedulableEvaluator.java、scheduler/…/TriggerEventHandler.java:343-344。

  5. 消息骨架:Queue 抽象与持久化 —— 把前四章串起来的"血管"。讲 QueueInterface 的 emit/receive 语义、QueueFactoryInterface 列出的所有队列种类,以及 JDBC 实现如何把队列做成一张表。这里能看清"单机 H2 vs 集群 Postgres/MySQL"到底换了什么。依据:core/…/queues/QueueInterface.java、queue/…/QueueFactoryInterface.java、queue-jdbc/…/JdbcQueueFactory.java。

  6. 扩展与接入:插件系统与 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.javaFlow(tasks / errors / triggers 字段)
一次运行 + 任务运行列表core/src/main/java/io/kestra/core/models/executions/Execution.javaExecution(id:71 / taskRunList:84 / state:105)
单个任务的一次运行core/src/main/java/io/kestra/core/models/executions/TaskRun.javaTaskRun(state:63 / attempts:56)
状态枚举 + 不可变状态机值core/src/main/java/io/kestra/core/models/flows/State.javaState.Type(:239)、State.withState(:73)
执行引擎主循环(状态机)executor/src/main/java/io/kestra/executor/ExecutorService.javaprocess(:135)、handleNext(:459)、handleWorkerTasks(:957)
执行引擎服务(订阅队列、写回)executor/src/main/java/io/kestra/executor/DefaultExecutor.javarun(:200)、toExecution(:578)
可执行任务的契约core/src/main/java/io/kestra/core/models/tasks/RunnableTask.javaRunnableTask.run(RunContext)
Worker 主体与启动worker/src/main/java/io/kestra/worker/AbstractWorker.javaAbstractWorker.start(:120)
Worker 轮询循环骨架worker/src/main/java/io/kestra/worker/WorkerLoop.javaWorkerLoop.runLoop(:88)
任务运行时环境core/src/main/java/io/kestra/core/runners/RunContext.javarender / logger / storage / metric
调度器主体scheduler/src/main/java/io/kestra/scheduler/DefaultScheduler.javaDefaultScheduler(:45)
触发器评估scheduler/src/main/java/io/kestra/scheduler/internals/SchedulableEvaluator.javaSchedulableEvaluator.evaluate(:36)
触发器→执行 的落地scheduler/src/main/java/io/kestra/scheduler/TriggerEventHandler.javatoExecution + triggerExecutionPublisher.send(:343-344)
触发器发布到命令队列core/src/main/java/io/kestra/core/scheduler/service/DefaultTriggerExecutionPublisher.javaDefaultTriggerExecutionPublisher(executionCommandQueue:26)
队列抽象(emit/receive)core/src/main/java/io/kestra/core/queues/QueueInterface.javaQueueInterface(emit:16 / receive:74)
所有队列种类清单queue/src/main/java/io/kestra/queue/QueueFactoryInterface.javaQueueFactoryInterface(:20)
JDBC 队列工厂queue-jdbc/src/main/java/io/kestra/queue/jdbc/JdbcQueueFactory.javaJdbcQueueFactory
Web 层:执行命令入口webserver/src/main/java/io/kestra/webserver/controllers/api/ExecutionController.javaexecutionCommandQueue.emit(:640)、webhook(:524)
插件注册表(type→类)core/src/main/java/io/kestra/core/plugins/PluginRegistry.javafindClassByIdentifier(:66)
插件扫描/加载core/src/main/java/io/kestra/core/plugins/PluginScanner.javaPluginScanner.scan(:59)
单机启动(H2 后端)cli/src/main/java/io/kestra/cli/commands/servers/LocalCommand.javaLocalCommand(h2 配置:33-40)
单 JVM 跑全部服务cli/src/main/java/io/kestra/cli/commands/servers/StandAloneCommand.javaStandAloneCommand(:32)

一句话回顾: YAML → Flow;点运行/触发 → ExecutionCommand 进队列 → Executor 无状态状态机推进 → Worker 跑任务回传 → 直到终态。所有连接线都是队列,而队列和仓储可换后端——这就是 Kestra 既能单机又能扛百万级的全部秘密。