领域模型:Flow / Task / Execution / State
30 秒导读: Kestra 是一个"写 YAML 就能编排任务"的工作流引擎。这一章只讲名词和数据模型——一份 YAML 怎么变成内存里的对象(
Flow),运行时又怎么被表示(Execution/TaskRun),以及所有运转状态都压进哪一台状态机枚举(State.Type)。只讲"东西是什么、状态有哪些",不讲"状态怎么推进"——那是第 2 章 执行引擎的活。
1. 这是什么(零基础也能懂)
一句话定义: Kestra 的领域模型,就是"一份工作流 YAML"在代码里对应的那几个 Java 类。
先建立最重要的一个直觉——蓝图 vs 实例:
| 你写的 | 代码里的类 | 类比 |
|---|---|---|
| 一份 flow YAML(定义) | Flow | 菜谱 / 类(class) |
| 点一次"运行"产生的一次运行 | Execution | 按菜谱做的这一顿饭 / 对象(instance) |
| flow 里的一个 task 定义 | Task(子类) | 菜谱里的一个步骤 |
| 那个 task 这次实际跑 出来的记录 | TaskRun | 这顿饭里"切菜"这一步的实况 |
| 任何东西此刻处于什么阶段 | State | 步骤旁边贴的"进行中/已完成"便签 |
一句话记牢:Flow / Task 是定义(写死在 YAML 里,不动);Execution / TaskRun 是运行时(每跑一次新生一批);State 贴在后两者身上记录它们走到哪一步了。
用起来什么样。 一份最小的 flow YAML 长这样:
id: hello
namespace: company.team
inputs:
- id: name
type: STRING
defaults: World
tasks:
- id: greet
type: io.kestra.plugin.core.log.Log
message: "Hello {{ inputs.name }}"
这份文本被解析后,id/namespace/inputs/tasks 分别落到 Flow 对象的对应字段上。点一次运行,引擎就为它造一个 Execution,并为 greet 这个 task 造一个 TaskRun,两者的 State 从 CREATED 开始往前走。
2. 顶层全景(这几个对象怎么串起来)
怎么读这张图: 左边是"定义层"(YAML 解析出来的,只读蓝图),右边是"运行时层"(每次运行新造的实例)。中间那道竖线就是"从蓝图实例化"这一步。
定义层 (蓝图, 来自 YAML) 运行时层 (每次运行新建)
┌──────────────────────────┐ ┌──────────────────────────┐
│ Flow │ │ Execution │
│ ├─ id / namespace │ 实例化 │ ├─ id / flowId │
│ ├─ inputs: List<Input> │ ──────▶│ ├─ state: State │
│ ├─ outputs: List<Output>│ │ └─ taskRunList ┐ │
│ └─ tasks: List<Task> ┐ │ └──────────────────┼───────┘
└────────────────────────┼──┘ │
│ ▼
每个 Task ┌──────────────────────────┐
┌──────────┴─────────┐ │ TaskRun │
│ RunnableTask (跑) │ ─────▶│ ├─ taskId (指回定义) │
│ FlowableTask (控) │ │ ├─ state: State │
│ ExecutableTask(派) │ │ └─ attempts: │
└────────────────────┘ │ List<TaskRunAttempt>│
└──────────────────────────┘
State (枚举状态机, 贴在 Execution / TaskRun / Attempt 上)
CREATED → RUNNING → SUCCESS / FAILED / KILLED / ...
各部件一句话职责:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
Flow | 一份 flow 的完整定义:任务列表、输入、输出、触发器 | core/models/flows/Flow.java |
AbstractFlow | Flow 的公共骨架:id、namespace、revision、inputs | core/models/flows/AbstractFlow.java |
Task | 单个任务定义的抽象基类;所有插件任务都继承它 | core/models/tasks/Task.java |
Execution | 一次运行的完整快照:所有 TaskRun + 整体 state | core/models/executions/Execution.java |
TaskRun | 一个 task 在这次运行里的实例 + 状态 + 尝试记录 | core/models/executions/TaskRun.java |
State | 状态机:当前状态 + 历史,判定"跑完没/跑着没" | core/models/flows/State.java |
YamlParser | YAML 文本 → 对象的入口 | core/serializers/YamlParser.java |
主线走一遍(高层): YAML 文本 → YamlParser.parse 反序列化成 Flow → 触发运 行时 Execution.newExecution(flow, ...) 造出 Execution → 引擎为每个待跑 task 造 TaskRun → 每个对象身上的 State 从 CREATED 逐步推进到终态。
3. 定义层:声明式 Flow
3.1 Flow 与它的继承链
一份 flow 的字段被拆到两层类里,靠 Lombok 的 @SuperBuilder 拼起来:
FlowInterface (契约)
▲
AbstractFlow ← 公共身份字段: id / namespace / revision / inputs / outputs / labels
▲
Flow ← 编排内容: tasks / errors / finally / triggers / outputs / concurrency
▲
FlowWithSource ← 多带一个字段: source (原始 YAML 文本)
AbstractFlow 放"身份"字段。 id、namespace、revision、inputs、outputs、disabled 都在这里(core/models/flows/AbstractFlow.java:33-56,AbstractFlow)。注意几个校验约束是写死在注解上的:
id必须匹配^[a-zA-Z0-9][a-zA-Z0-9._-]*,长度 1–100(AbstractFlow.java:31-33)。namespace只能小写,长度 1–150(AbstractFlow.java:36-38)。revision(修订号)@Min(1),同一个 flow 每次改动 +1(AbstractFlow.java:40-41)。
Flow 放"编排内容"字段。 真正描述"跑什么、怎么跑"的都在这层(core/models/flows/Flow.java:48,Flow extends AbstractFlow):
| 字段 | 类型 | 含义 | 行号 |
|---|---|---|---|
tasks | List<Task> | 主任务列表,@NotEmpty | Flow.java:72 |
errors | List<Task> | 出错时跑的任务 | Flow.java:75 |
_finally | List<Task> | 无论成败都跑(YAML 里叫 finally) | Flow.java:80 |
afterExecution | List<Task> | 执行结束后的钩子任务 | Flow.java:87 |
triggers | List<AbstractTrigger> | 触发器(定时/事件) | Flow.java:90 |
outputs | List<Output> | 暴露给其它 flow 的输出值 | Flow.java:104 |
concurrency | Concurrency | 并发上限策略 | Flow.java:96 |
finally是 Java 关键字,不能当字段名,所以内部字段叫_finally,再手写一个getFinally()把它以正常名字暴露出去(Flow.java:78-84)。这是个纯粹为绕开语言关键字的小技巧。
Flow 还提供一批"遍历任务树"的工具方法,后面章节会反复用到:
allTasks():把 tasks + errors + finally + afterExecution 平铺成一条流(Flow.java:138)。allTasksWithChilds():递归展开——遇到 flowable 任务(见 §4)就下钻它的子任务(Flow.java:148-169)。findTaskByTaskId(id):按 id 找定义,找不到抛InternalException(Flow.java:222)。
3.2 FlowWithSource:为什么要多一个"带源码"的版本
Flow 刻意不保留原始 YAML 文本——它的 getSource() 直接返回 null,注释说得很直白:"保守起见,flow 绝不返回任何 source"(Flow.java:300-304)。
真正要展示/存原文时用子类 FlowWithSource,它就多一个 source 字段并 override getSource() 把文本吐出来(core/models/flows/FlowWithSource.java:15、:45)。两个方向的转换都手写了字段拷贝:toFlow() 剥掉 source 降级成 Flow(FlowWithSource.java:17),静态 of(flow, source) 反过来包上 source(FlowWithSource.java:57)。
类头注释说
Flow计划被废弃、推荐用FlowWithSource(Flow.java:37-41)。当前代码里两者仍并存。
3.3 Input:一个 type 字段分出 18 种子类
Input<T> 是抽象类,靠 Jackson 的 @JsonTypeInfo + @JsonSubTypes,用 YAML 里的 type: 字段决定反序列化成哪个具体子类(core/models/flows/Input.java:29-51)。共 18 种,例如:
type: STRING → StringInput
type: INT → IntInput
type: FILE → FileInput
type: SELECT → SelectInput
type: SECRET → SecretInput
type: FORM → FormInput (可嵌子输入)
...共 18 种,见 Input.java:32-49
每个 Input 带 id / type / required(默认 true)/ defaults / prefill 等字段(Input.java:60-91),并要求子类实现 validate(T input)(Input.java:98)。
一个不显然的精华:FORM 类型可以包一层子输入,解析后要"拍平"。expandToLeaves() 把每个 FormInput 的孩子复制成一份 id 被改写成点号路径(environment + . + region → environment.region)的叶子输入(Input.java:112-130)。复制手段也很特别——走一次 Jackson round-trip(toMap 再 toMap 回来),因为这个抽象类没有 toBuilder(),而 round-trip 能顺带把子类类型重新解析回来(Input.java:140-144,copyWithId)。
3.4 Output:flow 对外暴露的返回值
Output 比 Input 简单得多:id / description / value(可为动态表达式)/ type / required(core/models/flows/Output.java:18)。它是 flow 级别的"返回值",给别的 flow 引用用(Flow.java:98-104 的字段说明)。
注意
outputs字段在AbstractFlow(:52)和Flow(:104)里都声明了——Flow这个覆盖版带了@PluginProperty(dynamic = true),支持动态表达式。
4. 任务体系:三大接口分野
这是整个领域模型里最需要讲清楚的一处。所有 task 的定义都继承抽象类 Task(core/models/tasks/Task.java:39,Task implements TaskInterface),它只提供公共字段(id、type、retry、timeout、disabled、runIf 等,Task.java:41-92)。
但一个 task "属于哪一类",不看它继承谁,而看它实现了下面哪个接口。 三个接口互斥地回答"这个任务由谁来跑、跑出来是什么":
| 接口 | 一句话 | 由谁执行 | 关键方法 | 文件 |
|---|---|---|---|---|
RunnableTask<T> | 真正干活的叶子任务 | Worker | run(RunContext) | RunnableTask.java:10 |
FlowableTask<T> | 控制流程,不干活 | Executor | childTasks / resolveNexts | FlowableTask.java:21 |
ExecutableTask<T> | 派生出子流程执行 | Executor | createSubflowExecutions | ExecutableTask.java:20 |
Task 用一个 isFlowable() 一票判定归属:return this instanceof FlowableTask(Task.java:154-157)。
4.1 RunnableTask —— 在 Worker 里跑
最简单。整个接口就一个方法:
// RunnableTask.java:10
public interface RunnableTask<T extends Output> extends Plugin, WorkerJobLifecycle {
// 在 Worker 里被调用来真正执行任务
T run(RunContext runContext) throws Exception;
}
一句话:凡是"下载文件、发 HTTP 请求、跑一段脚本、写日志"这种真正产生副作用的任务,都是 RunnableTask,由 Worker 进程调 run() 执行(详见第 3 章 Worker 与 RunContext)。返回值 T extends Output 是这次运行产出的结构化输出。
4.2 FlowableTask —— 控制流程,自己不干活
Parallel、Switch、EachSequential、Sequential 这些"控制流"任务实现的是 FlowableTask。它们自己不产生副作用,只决定"接下来跑哪些子任务"。所以它们由 Executor 在编排循环里处理,而不是发给 Worker。
三个承重方法:
allChildTasks()—— 返回全部子任务(含 errors)(FlowableTask.java:45)。childTasks(runContext, parentTaskRun)—— 解析出这一层要跑的子任务;对迭代型(如EachSequential)会把所有迭代都展开(FlowableTask.java:53)。resolveNexts(runContext, execution, parentTaskRun)—— 决定"下一步跑哪些",返回List<NextTaskRun>;串行的返回一个,并行的返回"并发数"那么多个(FlowableTask.java:61)。
它还提供一个默认的 resolveState(...),把子任务们的状态汇总成父任务的状态(FlowableTask.java:76-87)——但具体怎么汇总、怎么推进是第 2 章的事,这里只需知道"父 flowable 的状态是子任务状态的函数"。
4.3 ExecutableTask —— 派生子流
ExecutableTask 专门给"