mobile wallpaper 1
mobile wallpaper 2
mobile wallpaper 3
mobile wallpaper 4
mobile wallpaper 5
mobile wallpaper 6
12053 字
32 分钟
工作流引擎学习
2026-09-15

工作流引擎学习笔记:从任务编排到可靠执行#

目录#


1. 先定义:你要做的到底是什么#

1.1 工作流、任务、编排、调度不是同一个概念#

概念解决的问题典型例子
任务(Task)一次可执行工作调用 HTTP 接口、运行脚本、发送消息
工作流(Workflow)多个任务的结构与控制逻辑A 成功后并行执行 B/C,二者完成后执行 D
编排(Orchestration)由中心控制器决定下一步引擎保存状态并主动驱动下游任务
协同/编舞(Choreography)各服务通过事件自行反应服务 A 发布事件,服务 B/C 分别订阅
调度(Scheduling)决定何时获得执行机会定时、延迟、优先级、限流、资源匹配
执行器(Worker/Executor)真正执行任务代码拉取任务、执行业务逻辑、上报结果
事务(Transaction)约束一组数据操作的原子性数据库事务、消息发送与状态更新

**关键判断:**工作流引擎不是“把几个 API 串起来”的脚本工具。它的核心价值是:当任务跨越网络、进程、机器、服务和很长时间后,仍然能够回答“现在处于哪一步、下一步是什么、失败后怎么办、重启后如何继续”。


2. 工作流引擎必须掌握的概念#

2.1 定义、版本、实例、运行、尝试#

建议严格区分以下对象:

  • Workflow Definition(工作流定义):描述节点、边、参数、策略和版本。
  • Workflow Definition Version(定义版本):发布后不可变的具体版本,例如 order-sync@3
  • Workflow Instance(工作流实例):一次业务运行,例如订单 O-1001 的同步流程。
  • Workflow Run/Execution(执行):实例在引擎中的一次运行;如果支持恢复或重跑,要明确它与实例的关系。
  • Task/Activity(任务/活动):定义中的一个可执行节点。
  • Task Attempt(任务尝试):同一任务因为重试而产生的第 1、2、3 次实际执行。
  • Worker(执行器):领取任务并执行的进程或服务。
  • Event(事件):推动状态变化的事实,如 TaskSucceededTimerFiredSignalReceived
  • Signal/Command(信号/命令):外部希望工作流发生变化的输入;信号是事实,命令是请求,语义不要混用。
  • Compensation(补偿):业务层面的反向动作,不等同于数据库回滚。

一个常见错误是只设计一张 workflow_instance 表,用 status 字段覆盖所有含义。这样很快会遇到:同一个任务重试次数无法解释、历史状态丢失、并发更新互相覆盖、恢复位置不明确、审计无法还原。

2.2 控制流与数据流#

工作流有两条同时存在的图:

  1. 控制流图:节点何时可执行,边如何连接,条件、并行、汇聚、循环如何表达。
  2. 数据流图:节点读取哪些输入、产生哪些输出、输出如何映射到下游。

例如:

输入 orderId
|
加载订单 ──失败──> 重试/失败终止
|
+──> 计算库存 ──+
| +──> 创建发货单
+──> 计算优惠 ──+

控制流解决“什么时候做”;数据流解决“拿什么做”。将所有数据都塞进一个可变的全局 JSON 上下文很简单,但会造成大对象复制、字段覆盖、隐式依赖和敏感信息扩散。更稳妥的方式是:上下文只保存小型编排数据,业务大对象放业务系统或对象存储,任务之间传递引用和版本。

2.3 同步、异步、等待与外部完成#

任务至少有三种生命周期:

READY -> RUNNING -> SUCCEEDED
-> FAILED
-> TIMED_OUT
-> CANCELED
-> RETRY_WAITING -> READY

还有一种重要状态:

RUNNING -> WAITING_EXTERNAL -> SUCCEEDED/FAILED/TIMED_OUT

它适合等待外部回调、消息、设备上报或异步作业完成。不要让 worker 长时间占用线程等待;应该持久化等待条件,由事件接收器或定时器唤醒工作流。

2.4 状态不是事实,事件才是事实#

状态是当前投影,事件是导致状态变化的事实:

WorkflowStarted
TaskScheduled
TaskDispatched
TaskStarted
TaskSucceeded
TimerCreated
TimerFired
WorkflowCompleted

事件至少要包含:事件 ID、工作流实例 ID、序列号、类型、发生时间、载荷摘要、来源、关联键和幂等键。是否完整采用事件溯源是架构选择,但即便使用当前状态表,也强烈建议保留不可变的执行历史。

2.5 失败类型#

把错误分成不同类别,才能制定正确策略:

  • 业务失败:输入不满足业务规则,通常不应盲目重试。
  • 瞬时基础设施失败:连接超时、限流、临时不可用,通常适合退避重试。
  • 永久失败:参数错误、资源不存在、权限不足,重试通常无效。
  • 执行器失败:worker 崩溃、租约过期、进程被杀,平台需要重新投递。
  • 未知结果:请求发出后客户端超时,但远端可能已经成功,这是最需要幂等和查询确认的情况。
  • 编排错误:定义非法、节点不存在、表达式错误,应在发布或启动前尽量发现。

不要只存一个 error_message。建议同时保存错误类别、可重试标记、错误码、原始错误摘要、堆栈引用、首次发生时间、最近发生时间和尝试次数。


3. 三类核心架构模型#

工作流引擎没有唯一架构,最重要的差异在“谁保存状态、谁决定下一步、任务如何被表达”。

3.1 数据库驱动的持久化状态机#

**模型:**把每个节点和实例存成数据库记录;调度器扫描可执行记录,领取后投递;任务回报结果,事务性更新状态。

API -> Definition Store
-> Instance Store -> Scheduler -> Queue -> Worker
^ |
+-------------+

优点:

  • 容易理解和自研,SQL 可查询;
  • 状态恢复路径清晰;
  • 适合中短流程、规则较简单的编排;
  • 可以较快实现暂停、重试、定时和管理后台。

缺点:

  • 扫描和锁竞争可能成为瓶颈;
  • 复杂分支、循环、动态子流程会让状态表复杂;
  • 如果状态更新与外部副作用没有设计好,容易重复执行;
  • 长时间等待需要可靠的定时器和唤醒机制。

**适用:**第一版通用任务编排、内部平台、任务数量中等、团队擅长关系数据库。

**关键设计:**使用版本号或乐观锁;任务领取采用租约;状态迁移必须校验前置状态;调度查询要有索引和分片策略。

3.2 事件溯源 + 确定性重放#

**模型:**工作流逻辑根据历史事件重放,得到当前内存状态;新事件追加到历史中;事件历史是恢复和审计的主要依据。

Event History -> Replay Workflow Logic -> Commands/Timers/Tasks
^ |
+---------- Append New Events <--------+

Temporal 的工作流模型、事件历史和确定性约束是这类架构的代表性资料,可阅读其工作流文档事件历史说明工作流定义约束

优点:

  • 长时间运行、暂停恢复、故障恢复能力强;
  • 历史天然适合审计、调试和时间线展示;
  • 代码可以表达复杂控制流;
  • 状态迁移不依赖频繁保存完整快照。

缺点:

  • 必须理解确定性重放、版本兼容和事件历史增长;
  • 外部副作用不能随意放在可重放逻辑中;
  • 运维和调试门槛高于普通状态机;
  • 历史、快照、归档和查询投影需要额外设计。

**适用:**持续时间长、流程复杂、可靠恢复比简单查询更重要的系统。

3.3 DAG 调度器#

**模型:**工作流主要是有向无环图,节点通常是批处理、数据管道或容器任务;调度器根据依赖关系启动节点。

extract -> transform -> load
\-> validate

Apache Airflow 将工作流定义为任务依赖图,可参考其核心概念总览;Dagster 强调资产和数据血缘,可参考概念文档;Prefect 的 flow/task 模型可参考Flows

**优点:**依赖关系直观;批处理、数据管道和定时运行体验好;历史运行、重试和任务日志通常较成熟。

**缺点:**传统 DAG 不擅长任意循环、事件等待、业务长事务和细粒度外部交互;频繁动态创建实例、秒级响应和高并发在线编排可能需要额外改造。

**适用:**数据工程、ETL、离线作业、定时批任务。

3.4 三类模型的选择结论#

需求首选模型原因
先做可控 MVP数据库状态机语义清晰、实现成本低
长时间运行、强恢复事件历史/持久化执行能从事件重建状态
ETL、批任务、数据资产DAG 调度器依赖与数据血缘更自然
复杂在线业务编排状态机 + 事件历史兼顾查询、恢复和事件输入
容器级批处理Kubernetes 工作流任务隔离和资源调度更重要

不要因为“工作流”三个字就直接选择 BPMN,也不要因为“任务”两个字就直接选择消息队列。先看持续时间、执行边界、恢复语义和数据一致性。


4. 工作流定义语言与建模方式#

4.1 JSON/YAML DSL#

示例:

id: order-sync
version: 1
input:
- orderId
steps:
- id: load-order
type: http
method: GET
endpoint: ${services.order}/orders/${input.orderId}
retry:
maxAttempts: 3
backoff: exponential
next: [parallel-work]
- id: parallel-work
type: parallel
branches: [calculate-stock, calculate-discount]
join: all
next: [create-shipment]
- id: calculate-stock
type: task
worker: inventory
- id: calculate-discount
type: task
worker: pricing
- id: create-shipment
type: task
worker: shipping
timeout: 30s

**优点:**易于版本化、导入导出和自动生成;适合平台化和多语言执行器。

**缺点:**表达式、类型系统、错误提示、IDE 支持和静态校验需要自己建设。

4.2 代码优先#

使用宿主语言定义流程,例如通过函数、条件、循环、并发原语描述工作流。

**优点:**复用语言能力,复杂逻辑自然,测试工具成熟。

**缺点:**需要处理确定性、版本兼容、沙箱和安全;跨语言调用和可视化较复杂。

4.3 BPMN 2.0#

BPMN 是 OMG 发布的业务流程建模标准,官方入口见BPMN 2.0.2 规范。它对事件、任务、网关、子流程、消息、定时器等有统一图形语义。

**优点:**标准化、图形表达强、业务沟通好、工具生态成熟。

**缺点:**规范范围大;很多语义偏向业务流程和人工协作;直接把所有技术执行细节塞进 BPMN 会造成图过于复杂。

本项目不涉及角色审批时,BPMN 仍可以作为参考语义,但不必一开始完整实现 BPMN。更简单的做法是先实现自己的有限 DSL,再提供 BPMN 子集导入或导出。

4.4 CNCF Serverless Workflow#

Serverless Workflow Specification提供了面向服务编排的声明式规范,包含状态、事件、分支、重试、错误处理等概念。它适合研究“可移植的声明式编排 DSL”,但采用标准不等于自动获得可靠执行能力,持久化、调度、worker 协议仍要落地。

4.5 DSL 设计建议#

第一版只实现少数稳定原语:

  • task:投递给 worker 的任务;
  • http:受控的 HTTP 调用;
  • parallel:并行分支和汇聚;
  • choice:条件分支;
  • sleep/timer:延迟和定时;
  • wait-event:等待外部事件;
  • subflow:调用另一个版本化工作流;
  • compensate:显式补偿动作;
  • map:有界集合并行处理,避免无限扩张。

每个节点至少拥有:idtype、输入映射、输出映射、超时、重试、错误策略、并发限制、资源/队列标签和版本兼容信息。

4.6 静态校验清单#

发布定义时应检查:

  • 节点 ID 唯一;
  • 起点和终点明确;
  • 所有引用的节点、变量、worker 类型存在;
  • 图结构没有意外死循环;
  • 并行分支有合法汇聚策略;
  • 条件表达式可解析、变量类型匹配;
  • 超时、重试和并发参数在合法范围;
  • 敏感配置只引用 secret 名称,不把密钥写进定义;
  • 定义可序列化且包含 schema 版本;
  • 发布后不可变,新修改必须产生新版本。

5. 执行语义:状态机、事件溯源与持久化#

5.1 一个可执行节点的最小状态机#

PENDING
-> READY
-> DISPATCHED
-> RUNNING
-> SUCCEEDED
-> FAILED
-> RETRY_WAITING
-> TIMED_OUT
-> CANCELED
-> UNKNOWN

推荐将 UNKNOWN 单独建模:当执行器或网络在副作用发生后失联时,引擎不能武断地把任务标记为失败并立即重复执行。它需要通过幂等键、查询接口、回调或人工修复确认结果。

状态迁移要满足:

  • 合法性:只允许定义过的迁移;
  • 原子性:状态与版本号/事件序列一起更新;
  • 单调性:不能因迟到的旧消息把成功改回运行;
  • 可重试性:重复收到同一结果不应破坏状态;
  • 可解释性:每次迁移可追溯到事件和触发者。

5.2 数据库状态机的关键并发控制#

一个典型领取任务的逻辑是:

  1. 找到 READYavailable_at <= now 的任务;
  2. 通过行锁、乐观锁或原子更新抢占;
  3. 写入 lease_ownerlease_untilattempt
  4. 提交事务后执行外部任务;
  5. 上报结果时校验租约或 fencing token;
  6. 结果写入只允许作用于当前尝试。

不要只依靠 status = RUNNING。如果 worker 崩溃,必须通过租约过期恢复;如果旧 worker 在网络分区后恢复,fencing token 可以阻止它覆盖新 worker 的结果。

5.3 事件历史和快照#

若使用事件溯源,可以按以下方式折中:

  • 事件表保存完整状态变化;
  • 每 N 个事件或达到大小阈值保存快照;
  • 恢复时读取最新快照,再重放后续事件;
  • 历史超过阈值后归档到对象存储;
  • 读查询使用投影表,不要每次都重放全历史。

事件 schema 必须版本化。新增字段通常可以向后兼容;改变旧字段含义或删除字段则需要 upcaster、迁移策略或新事件类型。

5.4 确定性与副作用边界#

在可重放架构中,工作流逻辑可以决定“应该执行什么”,但不能在重放时再次直接执行外部副作用。以下操作通常必须放到 activity/task 中:

  • 读写数据库;
  • 调用 HTTP、RPC、消息系统;
  • 读取当前时间或随机数;
  • 访问文件系统、环境变量或外部配置;
  • 发送邮件、扣款、创建资源。

工作流层应获得一个稳定的结果,并将结果作为历史的一部分保存。这样恢复时只重放决策,不重复产生副作用。

5.5 至少一次、至多一次与有效一次#

分布式系统中,“只执行一次”通常不是单靠消息投递就能保证的。应把语义拆开:

  • 至少一次投递:失败会重新发送,可能重复;
  • 至多一次投递:发送后不重试,可能丢失;
  • 业务有效一次:允许底层重复,但通过幂等键、唯一约束或去重使业务效果只生效一次。

自研引擎的默认目标应是:任务投递至少一次,任务执行器和业务 API 通过幂等键实现有效一次。


6. 分布式可靠性:重试、幂等、超时与补偿#

6.1 重试策略#

常见策略:

  • 固定间隔:简单,但容易形成重试风暴;
  • 线性退避:间隔逐渐增加;
  • 指数退避:适合瞬时故障;
  • 指数退避 + 抖动:避免大量实例同时重试;
  • 按错误码、节点类型和业务状态决定是否重试。

一个常用形式:

delay = min(maxDelay, baseDelay * 2^(attempt - 1)) + random(0, jitter)

重试配置至少包含最大尝试次数、最大总时长、可重试错误集合、退避策略、抖动、超时和最终失败动作。重试不是可靠性的全部;如果副作用不可幂等,重试可能使问题更严重。

6.2 幂等设计#

幂等键通常由业务语义决定,例如:

idempotency_key = workflow_instance_id + task_id + attempt_semantic_key

但注意:如果每次 retry 都把 attempt 编号放入幂等键,重试就失去去重作用。通常要区分:

  • 任务逻辑幂等键:同一个业务动作使用同一个 key;
  • 投递消息 ID:每次投递唯一,用于追踪;
  • 尝试 ID:区分实际执行尝试。

服务端应通过唯一索引、幂等记录表或业务状态机记录已处理请求。执行器回报结果也必须幂等:同一 attempt_id 的成功回报重复到达时只接受一次。

6.3 超时不是取消#

超时只表示“引擎不再等待这次尝试”,不一定表示远端工作停止。必须分别设计:

  • 连接超时:无法建立连接;
  • 执行超时:任务运行超过限制;
  • 心跳超时:worker 不再证明自己存活;
  • 工作流超时:整个实例超过截止时间;
  • 取消:向任务和下游发送停止请求;
  • 截止时间(deadline):所有重试、排队和执行共享的绝对时间边界。

6.4 Saga 与补偿#

跨多个服务的业务动作通常无法用一个数据库事务回滚。Saga 将长事务拆成多个本地事务,并为已完成动作定义补偿动作。经典论文可参考 Garcia-Molina 与 Salem 的Sagas

示例:

冻结库存 -> 创建支付单 -> 创建发货单
| | |
解冻库存 <- 取消支付单 <- 取消发货单

补偿不是严格逆操作,也不保证回到最初状态。例如已经发送的邮件不能真正撤回,只能发送更正通知。补偿动作同样需要幂等、超时、重试和审计。

6.5 资源锁与并发策略#

同一业务实体可能被多个工作流实例同时修改。可选策略:

  • 业务层乐观锁:带版本号更新;
  • 引擎级互斥键:同一 key 同时只允许一个节点运行;
  • 队列分区:同一 key 路由到同一分区;
  • 数据库唯一约束;
  • 顺序化子流程;
  • 分布式锁,但必须考虑租约、续期和 fencing。

优先使用业务版本号、唯一约束和分区顺序,而不是到处加分布式锁。


7. 调度、队列与并发控制#

7.1 调度器的职责#

调度器不是简单的定时扫描器,它至少负责:

  • 找出满足依赖的 READY 节点;
  • 判断定时器和截止时间;
  • 应用租户、队列、优先级和资源配额;
  • 控制工作流级、节点级和租户级并发;
  • 防止同一任务重复领取;
  • 处理租约过期和死 worker;
  • 让长期饥饿的任务得到机会;
  • 记录调度延迟和原因。

7.2 拉模式与推模式#

模式说明优点风险
Worker 拉取worker 主动领取任务背压自然、部署简单空轮询、长轮询实现复杂
Scheduler 推送调度器主动发送低延迟、中心控制强调度器和队列压力集中
队列中介调度器写队列,worker 消费解耦、可扩展需要处理重复、顺序、可见性超时
混合模式控制面调度,数据面拉取平衡吞吐和控制组件更多

MVP 建议使用“数据库记录 + 持久化队列/通知 + worker 拉取”的混合方式:数据库是事实来源,队列用于唤醒和削峰,队列丢失时可由扫描器兜底。

7.3 队列语义#

需要明确:

  • 是否保证顺序;
  • 是否至少一次;
  • 可见性超时如何设置;
  • 消费失败进入重试队列还是死信队列;
  • 消息确认在任务完成前还是领取后;
  • 消息是否包含完整输入还是只包含任务 ID;
  • 队列积压是否触发限流或告警。

工作流引擎不要把队列消息当成唯一状态。消息是通知,数据库或事件历史才是可恢复的事实来源。

7.4 定时器设计#

不要为每个未来任务启动一个线程。可采用:

  • 数据库 available_at 索引 + 周期扫描;
  • 时间轮;
  • 最小堆 + 持久化恢复;
  • 外部延迟队列;
  • 分层时间桶。

如果定时器是工作流语义的一部分,必须持久化:目标实例、目标节点、触发时间、截止策略、时区/日历、重复规则、状态和幂等标识。时区和夏令时必须在定义层面明确,避免只存本地时间。

7.5 并发与背压#

至少分三层限额:

  1. 租户级:防止一个租户占满平台;
  2. 工作流级:限制同一流程实例的并行分支;
  3. worker/资源级:限制 CPU、连接池、下游 API 配额。

背压信号包括队列深度、最老任务年龄、调度延迟、worker 利用率和下游限流率。扩容 worker 不能解决下游系统已经过载的问题。


8. 事件驱动与外部系统集成#

8.1 编排与事件驱动的关系#

事件适合表达“已经发生的事实”,命令适合表达“请执行某动作”。例如:

  • 事件:PaymentCaptured
  • 命令:CreateShipment
  • 查询:GetPaymentStatus

工作流可以订阅事件来推进等待状态,也可以发布命令驱动任务。不要把每个内部状态变化都暴露成公共事件;公共事件需要稳定 schema、版本和兼容策略。

8.2 CloudEvents#

CloudEvents 规范提供了事件元数据的通用格式,包括事件 ID、来源、类型、时间和内容类型。采用类似的信封有助于跨消息系统传递,但它不替你解决顺序、投递、幂等和事务一致性。

8.3 Outbox 与 Inbox#

当业务事务要同时写数据库和发送事件时,直接“先写库再发消息”或“先发消息再写库”都可能产生不一致。Outbox 的基本思路是:在同一个数据库事务中写业务数据和 outbox 记录,再由发布器异步发送;消费方使用 inbox/去重表避免重复处理。

一个可靠链路通常是:

业务事务:业务表 + outbox
|
outbox publisher
|
消息系统
|
workflow event receiver
|
event deduplication
|
append event / advance state

注意:outbox 解决的是本地数据库与消息发布之间的一致性,不是所有跨系统事务的万能方案。

8.4 Webhook 与回调#

外部回调接口至少需要:

  • 签名验证;
  • 请求 ID 和事件 ID;
  • 时间戳和重放防护;
  • 幂等去重;
  • 快速确认和异步处理;
  • 事件 schema 版本;
  • 未知实例、迟到事件和重复事件的处理策略。

回调处理器不应直接在 HTTP 请求线程中执行整个工作流推进;应该先验证、持久化事件,再异步推进。

8.5 轮询与回调的比较#

方式优点缺点
轮询对外部系统要求低、实现稳定延迟和请求量较高
Webhook低延迟、节省请求需要公网入口、签名和重放防护
消息订阅吞吐高、解耦好需要事件基础设施和 schema 治理
混合回调优先、轮询兜底设计和运维更复杂

9. 数据模型与存储设计#

9.1 推荐的逻辑实体#

workflow_definition
workflow_definition_version
workflow_instance
workflow_node_instance
node_attempt
workflow_event
timer
external_signal
task_queue
worker_lease
artifact/reference

9.2 最小字段建议#

workflow_definition_version

  • definition_idversionschema_version
  • status:draft/published/retired;
  • definition_json
  • input_schemaoutput_schema
  • created_bycreated_at
  • 内容 hash,保证版本不可悄悄修改。

workflow_instance

  • instance_iddefinition_iddefinition_version
  • business_keytenant_id
  • statuscurrent_revision
  • input_refoutput_ref
  • started_atdeadline_atcompleted_at
  • parent_instance_idcorrelation_id

node_attempt

  • attempt_idinstance_idnode_id
  • attempt_nostatus
  • queue_nameworker_idlease_untilfencing_token
  • idempotency_key
  • available_atstarted_atfinished_at
  • input_refoutput_referror_ref

workflow_event

  • event_idinstance_idsequence
  • event_typepayloadpayload_schema_version
  • occurred_atrecorded_atsource
  • dedup_key
  • (instance_id, sequence) 和必要的 dedup_key 建唯一约束。

9.3 输入输出存储#

小输入、小输出可以存数据库 JSON;大对象、文件和日志应存对象存储,并在工作流中保存引用、hash、大小、媒体类型和保留期限。不要把密码、访问令牌和完整个人敏感数据无差别写入事件历史或日志。

9.4 一致性边界#

建议把数据库作为控制面的权威来源:

  • 状态推进、租约、事件追加在同一存储中保持原子;
  • 队列只承担投递和唤醒;
  • worker 结果必须带实例、节点和尝试身份;
  • 读模型可以异步最终一致;
  • 关键管理操作要记录审计事件。

9.5 分片和归档#

初期可按 tenant_idinstance_id 做逻辑分区;规模上来后再按时间、租户和实例 hash 组合分片。历史执行数据应设置保留期、冷热分层和归档策略,否则查询索引和存储成本会持续增长。


10. 可观测性、运维与安全#

10.1 三类可观测信号#

采用 OpenTelemetry Observability Primer中的三大信号思路:

  • Metrics:吞吐、成功率、重试率、队列深度、调度延迟、任务耗时、租约过期、定时器延迟;
  • Logs:结构化记录实例、节点、尝试、trace、租户和错误码;
  • Traces:从 API 请求到工作流实例、任务投递、worker 执行、外部调用的链路。

工作流特有的指标包括:

workflow_started_total
workflow_completed_total
workflow_failed_total
workflow_duration_seconds
workflow_queue_delay_seconds
node_attempt_total
node_retry_total
node_timeout_total
orphaned_attempt_total
stuck_instance_total

10.2 控制台必须回答的问题#

  • 这个实例当前在哪个节点?
  • 它为什么没有继续?是在等待依赖、队列、定时器还是外部事件?
  • 最近一次尝试由哪个 worker 执行?租约何时过期?
  • 失败是业务错误、基础设施错误还是未知结果?
  • 重试还剩多少次?下一次何时执行?
  • 哪些下游服务导致整体变慢?
  • 暂停、恢复、取消、重跑、跳过的影响是什么?

10.3 管理操作的安全边界#

暂停、取消、重试、跳过节点、修改输入、人工标记成功都属于高风险操作。即使第一版不做角色审批,也至少需要:

  • 明确的管理 API;
  • 操作者身份和原因;
  • 审计记录;
  • 防止绕过状态机;
  • 对重跑和补偿提供预览;
  • 重要操作二次确认或双人机制作为未来扩展。

10.4 多租户安全#

如果未来服务多个团队,设计时就引入 tenant_id

  • 所有实例、任务、事件、日志和对象引用带租户边界;
  • 查询必须强制带租户过滤;
  • worker 队列按租户或资源等级隔离;
  • 限制单租户实例数、并发数、历史量和事件大小;
  • Secret 通过外部密钥管理系统引用;
  • 任务代码不应获得引擎数据库的任意写权限。

11. 参考产品与架构对比#

以下对比用于学习架构,不等于直接推荐。请以各项目当前官方文档为准。

11.1 Temporal / Cadence 类持久化执行引擎#

  • 资料:Temporal WorkflowsActivitiesWorkflow Definition
  • 核心思路:持久化工作流执行、事件历史、任务队列、worker 和活动边界。
  • 擅长:长流程、可靠恢复、异步等待、代码优先工作流。
  • 学习重点:确定性重放、活动幂等、事件历史增长、版本兼容。
  • 不要照搬:如果你的流程只有几秒、只在单库内运行,完整事件历史可能过重。

11.2 AWS Step Functions#

  • 资料:Amazon States Language
  • 核心思路:声明式状态机,由托管服务负责执行、重试、分支、等待和服务集成。
  • 擅长:云服务编排、可视化状态机、托管运维。
  • 学习重点:状态机 DSL、集成模式、错误处理、服务边界。
  • 局限:强绑定云生态;跨云或私有化场景需要重新评估。

11.3 Azure Durable Functions#

  • 资料:Durable Functions overview
  • 核心思路:编排器函数、活动函数、实体函数和持久化执行。
  • 擅长:Serverless 场景中的长时间编排、定时器、扇出/扇入。
  • 学习重点:编排器确定性、重放、活动隔离和托管运行时。

11.4 Argo Workflows#

  • 资料:Argo Workflows 官方文档
  • 核心思路:Kubernetes CRD 驱动的容器原生工作流。
  • 擅长:容器任务、批处理、机器学习流水线、Kubernetes 资源生态。
  • 学习重点:声明式资源、容器隔离、artifact、DAG/steps、Kubernetes 调度。
  • 局限:若任务主要是低延迟 RPC 或长时间等待外部业务事件,容器工作流未必是最自然的模型。

11.5 Airflow#

  • 资料:核心概念
  • 核心思路:以 DAG 为中心的批处理和数据管道调度。
  • 擅长:定时数据作业、依赖管理、任务日志、数据团队协作。
  • 学习重点:调度周期、任务实例、重试、资源池、回填和数据管道可观测性。
  • 不宜直接当作:高频在线业务事务编排或复杂外部回调引擎。

11.6 Dagster / Prefect#

  • 资料:Dagster ConceptsPrefect Flows
  • Dagster 更强调软件定义资产、数据血缘、数据质量和可观测性。
  • Prefect 更强调 Python flow/task、动态流程和开发者体验。
  • 学习重点:任务抽象、动态 DAG、调度、重试、数据资产与运行元数据。

11.7 Netflix Conductor / Orkes Conductor#

  • 资料:Conductor OSS Documentation
  • 核心思路:声明式 JSON 工作流、任务队列、worker、决策任务和运行时服务。
  • 擅长:微服务编排、任务队列驱动和可视化管理。
  • 学习重点:任务状态、轮询 worker、系统任务、动态工作流和持久化。

11.8 Camunda / BPMN 引擎#

  • 资料:Camunda 文档入口
  • 核心思路:BPMN/DMN 等标准模型与流程运行时结合。
  • 擅长:标准化业务流程、事件、消息、人工任务和流程治理。
  • 学习重点:BPMN 执行语义、边界事件、消息关联、版本迁移。
  • 本项目取舍:不做角色审批时,可以学习其事件、定时器和流程版本思想,不必先实现完整 BPMN。

11.9 选型矩阵#

产品/模型DSL/代码长流程DAG/批处理容器原生事件等待可移植性学习价值
Temporal代码优先持久化执行
Step FunctionsDSL托管状态机
Durable Functions代码优先重放与编排器
ArgoYAML/DSL容器工作流
AirflowPython/DAG弱~中数据调度
DagsterPython/资产数据资产
PrefectPython动态 flow
ConductorJSON/worker微服务编排
CamundaBPMN/DSL标准流程语义
自研状态机自定义取决于实现取决于实现取决于实现取决于实现取决于实现贴合业务

12. 自研架构建议#

12.1 推荐的分层#

┌─────────────────────────┐
│ API / CLI / Admin UI │
└────────────┬────────────┘
┌────────────▼────────────┐
│ Control Plane │
│ Definition / Instance │
│ Validation / Command │
└────────────┬────────────┘
┌────────────▼────────────┐
│ Orchestrator │
│ State Transition │
│ Dependency Evaluation │
│ Timer / Signal Handling │
└───────┬─────────┬───────┘
│ │
┌──────────▼───┐ ┌──▼──────────┐
│ Durable Store │ │ Queue/Bus │
│ DB + Events │ │ Dispatch │
└──────────┬───┘ └──┬──────────┘
│ │
┌───────▼─────────▼───────┐
│ Worker SDK / Executors │
│ HTTP / Script / Container│
└─────────────────────────┘

12.2 控制面与数据面#

  • 控制面:定义发布、实例管理、调度决策、状态、事件、定时器、重试策略、管理操作。
  • 数据面:任务真正运行的 worker、容器、脚本、网络调用和资源隔离。

两者分开能降低引擎核心的安全风险:业务任务不应直接修改控制面数据库;worker 只通过 SDK/API 拉取任务和回报结果。

12.3 第一版建议的技术路线#

阶段 A:单进程、单数据库验证语义

  • JSON/YAML DSL;
  • 任务节点、顺序、条件、并行;
  • 实例和节点状态表;
  • 内存 worker 或本地 HTTP worker;
  • 状态迁移单元测试。

阶段 B:可靠异步执行

  • 持久化任务队列;
  • worker 拉取与租约;
  • 重试、超时、心跳;
  • 幂等键、fencing token;
  • 定时器和外部事件;
  • 事件历史与查询接口。

阶段 C:生产化控制面

  • 多租户、配额、权限、审计;
  • 指标、日志、trace;
  • 归档、清理、备份和恢复;
  • 版本兼容、迁移和灰度发布;
  • 故障注入和压测。

阶段 D:高级能力

  • 子流程、动态 map、补偿;
  • 代码优先 SDK 或 BPMN/Serverless Workflow 导入;
  • 容器执行、资源调度;
  • 多区域和灾备;
  • 人工任务、规则和 AI 任务扩展。

12.4 关键接口草案#

POST /workflow-definitions
POST /workflow-definitions/{id}/versions/{version}/publish
POST /workflow-instances
GET /workflow-instances/{instanceId}
POST /workflow-instances/{instanceId}/signal
POST /workflow-instances/{instanceId}/pause
POST /workflow-instances/{instanceId}/resume
POST /workflow-instances/{instanceId}/cancel
GET /workflow-instances/{instanceId}/events
GET /tasks/poll
POST /tasks/{attemptId}/heartbeat
POST /tasks/{attemptId}/complete
POST /tasks/{attemptId}/fail

接口语义要明确:哪些操作是幂等的、哪些会创建新版本、迟到的回报如何处理、取消是否强制、查询是否最终一致。

12.5 任务协议草案#

worker 领取任务时获得:

{
"attemptId": "a-123",
"instanceId": "i-456",
"definitionVersion": "order-sync@1",
"nodeId": "create-shipment",
"taskType": "shipping",
"idempotencyKey": "i-456:create-shipment",
"input": {"orderId": "O-1001"},
"deadlineAt": "2026-09-22T12:00:00Z",
"leaseUntil": "2026-09-22T11:05:00Z"
}

完成时必须带:

{
"attemptId": "a-123",
"result": {"shipmentId": "S-1"},
"status": "SUCCEEDED",
"workerId": "worker-7",
"fencingToken": 42
}

引擎拒绝实例、节点、尝试或 fencing token 不匹配的回报;重复回报返回已知结果,而不是再次推进流程。

12.6 为什么不建议一开始做微服务拆分#

自研早期最难的是状态语义,不是服务数量。把定义服务、实例服务、调度服务、定时服务、事件服务、查询服务一开始全部拆开,会增加跨服务一致性和调试成本。建议先做模块化单体:控制面边界清晰、数据事务集中、接口可替换;当吞吐、团队边界或故障隔离真正需要时再拆分。


13. 最小可行产品设计#

13.1 MVP 必须支持#

  • 定义发布和不可变版本;
  • 顺序任务、条件分支、并行汇聚;
  • worker 拉取、租约、心跳、完成和失败;
  • 指数退避重试和任务超时;
  • 工作流暂停、恢复、取消;
  • 定时等待和外部信号;
  • 实例状态、节点状态和事件时间线;
  • 基本幂等和重复回报处理;
  • 结构化日志和核心指标。

13.2 MVP 暂不支持#

  • 任意代码在引擎进程内执行;
  • 任意动态修改运行中定义;
  • 无边界的递归和动态并行;
  • 跨系统强事务;
  • 完整 BPMN;
  • 自动修复所有未知结果;
  • 多区域主动-主动;
  • 角色审批和复杂组织权限。

13.3 验收场景#

至少实现并验证以下流程:

  1. 顺序成功:A → B → C;
  2. 条件分支:根据输入选择 B 或 C;
  3. 并行汇聚:B/C 并行,全部成功后 D;
  4. 单任务瞬时失败:退避后成功;
  5. 永久失败:达到策略后流程失败;
  6. worker 崩溃:租约过期后重新投递;
  7. 迟到回报:旧尝试不能覆盖新尝试;
  8. 重复回报:不会重复推进后续节点;
  9. 外部事件:流程等待回调后继续;
  10. 重启恢复:引擎重启后从持久化状态继续;
  11. 取消:取消后不可再启动新的普通节点;
  12. 超时:任务超时与工作流总超时行为符合定义。

13.4 最小可用数据表#

workflow_definition_versions
workflow_instances
node_instances
node_attempts
workflow_events
timers
external_signals

可以先不引入单独的队列表,使用 node_attemptsready_at 和 worker 拉取接口实现;当吞吐或跨语言 worker 需求出现后,再引入消息系统。


14. 测试、故障注入与正确性证明#

14.1 单元测试#

  • DSL 解析和 schema 校验;
  • 图遍历、条件表达式和并行汇聚;
  • 每种状态迁移的合法性;
  • 重试退避计算;
  • 事件去重;
  • 输入输出映射;
  • 版本兼容和事件升级。

14.2 属性测试与模型测试#

把工作流状态机写成小模型,验证任意事件序列不会违反不变量:

  • 已完成节点不会再次变成待执行;
  • 同一尝试最多产生一个最终结果;
  • 汇聚节点不会在所有前置未满足时执行;
  • 已取消实例不会生成普通新任务;
  • 事件序列号严格递增;
  • 结果只能被对应的实例、节点和尝试接受。

TLA+ 官方入口适合学习如何用形式化模型描述并发状态和不变量;不一定要在第一版引入完整形式化验证,但可以借鉴这种思路写状态迁移模型。

14.3 集成测试#

  • 数据库事务回滚;
  • 队列重复投递和消息丢失后的扫描兜底;
  • worker 重启、网络断开、心跳延迟;
  • 任务完成后引擎进程崩溃;
  • 数据库主从切换或连接池耗尽;
  • 外部事件重复、乱序和迟到;
  • 时钟偏移和夏令时边界。

14.4 故障注入#

至少注入:

  • 在“副作用已发生、结果未提交”时杀死 worker;
  • 在状态更新提交前后杀死 orchestrator;
  • 让旧租约 worker 延迟回报;
  • 让队列重复投递同一任务;
  • 让下游接口返回 429、500、超时和成功后断连接;
  • 让定时器扫描器停止一段时间;
  • 让事件接收器重复消费同一事件。

14.5 压测指标#

  • 实例启动吞吐;
  • 单实例节点推进延迟;
  • 调度吞吐和 p99 延迟;
  • 队列积压恢复时间;
  • 数据库锁等待和写放大;
  • 事件历史增长速度;
  • worker 扩展后的吞吐曲线;
  • 单租户高负载对其他租户的影响。

14.6 使用 Raft 等资料的正确方式#

Raft 论文适合学习共识、日志复制和领导者选举,但不要因为引擎需要可靠状态就直接自研 Raft。优先使用成熟数据库、消息系统和托管协调服务;只有在明确需要自建一致性存储且团队具备运维能力时,才把共识算法纳入设计范围。


15. 学习路线与调研方法#

15.1 第一阶段:建立基础模型(1~2 周)#

学习:

  • 有限状态机、DAG、事件驱动;
  • 数据库事务、隔离级别、乐观锁;
  • 消息投递、可见性超时、死信;
  • 幂等、重试、超时和 Saga;
  • 分布式系统中的故障、时钟和网络分区。

输出:画出一个“订单同步”工作流的状态图、事件图和数据库实体图。

15.2 第二阶段:对比实现(2~3 周)#

建议至少动手:

  1. 用关系数据库写一个顺序/条件/并行状态机;
  2. 用队列接入 worker,加入租约和重复投递;
  3. 加入持久化定时器和外部信号;
  4. 阅读 Temporal 的 workflow/activity 与事件历史;
  5. 阅读 Step Functions 的状态语言;
  6. 阅读 Airflow 或 Dagster 的 DAG/资产模型;
  7. 在 Argo 中运行一个容器 DAG,观察 Kubernetes 资源边界。

输出:一张“同一流程在不同模型中的表达”对比表,以及一份取舍记录。

15.3 第三阶段:实现可用 MVP(4~8 周)#

顺序建议:

  1. DSL 和静态校验;
  2. 定义版本和实例创建;
  3. 状态迁移和事件记录;
  4. worker 协议和租约;
  5. 重试、超时、心跳;
  6. 定时器和外部信号;
  7. 控制台查询和操作;
  8. 故障注入与压测。

每加一种能力,都先写状态迁移测试和失败场景测试,不要先堆管理页面。

15.4 第四阶段:生产化(持续迭代)#

  • 明确 SLO:调度延迟、恢复时间、数据保留和可用性;
  • 做容量模型:实例数、节点数、事件数、日志量;
  • 做灾备演练和恢复演练;
  • 做版本迁移、灰度和回滚;
  • 引入多租户、配额、审计和密钥管理;
  • 建立工作流定义评审和发布规范。

15.5 调研每个产品时都问同样的问题#

  1. 工作流状态保存在哪里?
  2. 任务投递是至少一次还是其他语义?
  3. worker 崩溃后如何恢复?
  4. 副作用如何保证幂等?
  5. 长时间等待如何实现?
  6. 版本发布后,运行中的实例如何处理?
  7. 事件历史是否完整保存?如何归档?
  8. 是否支持暂停、取消、重跑和人工修复?
  9. 调度器如何分片和扩展?
  10. 数据库、队列和 worker 之间的事实边界是什么?
  11. 任务代码运行在哪个信任边界?
  12. 出现未知结果时系统如何避免重复副作用?

16. 结论与决策清单#

16.1 对当前项目的推荐结论#

  1. 模型:从数据库驱动的持久化状态机开始,同时保留不可变事件历史;
  2. 流程表达:先用小型 JSON/YAML DSL,不要一开始实现完整 BPMN;
  3. 执行方式:worker 拉取 + 租约 + 心跳 + fencing token;
  4. 可靠性:默认至少一次投递,业务有效一次由幂等键和唯一约束实现;
  5. 等待:把定时器、外部事件和回调持久化,禁止用线程长期阻塞;
  6. 数据:控制面数据库是事实来源,队列只负责投递和唤醒;
  7. 扩展:通过任务类型、worker SDK、子流程和补偿扩展,而不是让引擎直接执行任意业务代码;
  8. 演进:当流程复杂度和运行时长增长,再引入事件重放、快照、归档或代码优先模型;
  9. 架构:先模块化单体,验证语义后再按吞吐和故障隔离拆分;
  10. 验证:把重启、重复投递、迟到回报、未知结果和外部事件乱序列为一等测试场景。

16.2 开始编码前必须写下的 ADR#

  • 为什么选择状态机而不是完整 BPMN?
  • 状态表和事件历史谁是权威?
  • 任务投递与业务有效一次如何定义?
  • worker 租约和 fencing 如何工作?
  • 运行中实例遇到新定义版本怎么办?
  • 取消是尽力而为还是强保证?
  • 超时后如何处理可能已经成功的外部副作用?
  • 大输入输出和敏感数据放在哪里?
  • 哪些管理操作允许跳过状态机?
  • 什么时候从模块化单体拆成服务?

16.3 最终判断标准#

一个工作流引擎是否设计得好,不看它能画出多漂亮的流程图,而看它在以下场景下是否仍然可解释、可恢复、可测试:

  • worker 在外部副作用完成后崩溃;
  • 消息重复、乱序、迟到或丢失;
  • 引擎重启、数据库短暂不可用;
  • 一个任务需要等待数天;
  • 工作流定义已经升级,但旧实例仍在运行;
  • 一个租户突然产生大量并行任务;
  • 运维人员需要取消、重试或修复一个异常实例。

如果这些问题能用明确的状态、事件、租约、幂等键和审计记录回答,才说明引擎具备可靠的基础。


17. 参考资料总表#

标准与规范#

工作流引擎与编排平台#

分布式系统与可靠性#

建议继续检索的主题#

  • workflow patterns、persistent execution、durable execution;
  • workflow versioning、event history、deterministic replay;
  • idempotency、deduplication、fencing token;
  • outbox/inbox、saga orchestration、compensation;
  • timer wheel、distributed scheduler、backpressure;
  • workflow testing、fault injection、model checking;
  • workflow observability、trace context propagation;
  • multi-tenant orchestration、quota、fair scheduling。
分享

如果这篇文章对你有帮助,欢迎分享给更多人!

工作流引擎学习
https://niushaoxiong.top/posts/工作流引擎学习笔记/
作者
一只捡星星的熊
发布于
2026-09-15
许可协议
CC BY-NC-SA 4.0

部分信息可能已经过时

目录