工作流引擎学习笔记:从任务编排到可靠执行
目录
- 1. 先定义:你要做的到底是什么
- 2. 工作流引擎必须掌握的概念
- 3. 三类核心架构模型
- 4. 工作流定义语言与建模方式
- 5. 执行语义:状态机、事件溯源与持久化
- 6. 分布式可靠性:重试、幂等、超时与补偿
- 7. 调度、队列与并发控制
- 8. 事件驱动与外部系统集成
- 9. 数据模型与存储设计
- 10. 可观测性、运维与安全
- 11. 参考产品与架构对比
- 12. 自研架构建议
- 13. 最小可行产品设计
- 14. 测试、故障注入与正确性证明
- 15. 学习路线与调研方法
- 16. 结论与决策清单
- 17. 参考资料总表
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(事件):推动状态变化的事实,如
TaskSucceeded、TimerFired、SignalReceived。 - Signal/Command(信号/命令):外部希望工作流发生变化的输入;信号是事实,命令是请求,语义不要混用。
- Compensation(补偿):业务层面的反向动作,不等同于数据库回滚。
一个常见错误是只设计一张 workflow_instance 表,用 status 字段覆盖所有含义。这样很快会遇到:同一个任务重试次数无法解释、历史状态丢失、并发更新互相覆盖、恢复位置不明确、审计无法还原。
2.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 状态不是事实,事件才是事实
状态是当前投影,事件是导致状态变化的事实:
WorkflowStartedTaskScheduledTaskDispatchedTaskStartedTaskSucceededTimerCreatedTimerFiredWorkflowCompleted事件至少要包含:事件 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 \-> validateApache Airflow 将工作流定义为任务依赖图,可参考其核心概念总览;Dagster 强调资产和数据血缘,可参考概念文档;Prefect 的 flow/task 模型可参考Flows。
**优点:**依赖关系直观;批处理、数据管道和定时运行体验好;历史运行、重试和任务日志通常较成熟。
**缺点:**传统 DAG 不擅长任意循环、事件等待、业务长事务和细粒度外部交互;频繁动态创建实例、秒级响应和高并发在线编排可能需要额外改造。
**适用:**数据工程、ETL、离线作业、定时批任务。
3.4 三类模型的选择结论
| 需求 | 首选模型 | 原因 |
|---|---|---|
| 先做可控 MVP | 数据库状态机 | 语义清晰、实现成本低 |
| 长时间运行、强恢复 | 事件历史/持久化执行 | 能从事件重建状态 |
| ETL、批任务、数据资产 | DAG 调度器 | 依赖与数据血缘更自然 |
| 复杂在线业务编排 | 状态机 + 事件历史 | 兼顾查询、恢复和事件输入 |
| 容器级批处理 | Kubernetes 工作流 | 任务隔离和资源调度更重要 |
不要因为“工作流”三个字就直接选择 BPMN,也不要因为“任务”两个字就直接选择消息队列。先看持续时间、执行边界、恢复语义和数据一致性。
4. 工作流定义语言与建模方式
4.1 JSON/YAML DSL
示例:
id: order-syncversion: 1input: - orderIdsteps: - 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:有界集合并行处理,避免无限扩张。
每个节点至少拥有:id、type、输入映射、输出映射、超时、重试、错误策略、并发限制、资源/队列标签和版本兼容信息。
4.6 静态校验清单
发布定义时应检查:
- 节点 ID 唯一;
- 起点和终点明确;
- 所有引用的节点、变量、worker 类型存在;
- 图结构没有意外死循环;
- 并行分支有合法汇聚策略;
- 条件表达式可解析、变量类型匹配;
- 超时、重试和并发参数在合法范围;
- 敏感配置只引用 secret 名称,不把密钥写进定义;
- 定义可序列化且包含 schema 版本;
- 发布后不可变,新修改必须产生新版本。
5. 执行语义:状态机、事件溯源与持久化
5.1 一个可执行节点的最小状态机
PENDING -> READY -> DISPATCHED -> RUNNING -> SUCCEEDED -> FAILED -> RETRY_WAITING -> TIMED_OUT -> CANCELED -> UNKNOWN推荐将 UNKNOWN 单独建模:当执行器或网络在副作用发生后失联时,引擎不能武断地把任务标记为失败并立即重复执行。它需要通过幂等键、查询接口、回调或人工修复确认结果。
状态迁移要满足:
- 合法性:只允许定义过的迁移;
- 原子性:状态与版本号/事件序列一起更新;
- 单调性:不能因迟到的旧消息把成功改回运行;
- 可重试性:重复收到同一结果不应破坏状态;
- 可解释性:每次迁移可追溯到事件和触发者。
5.2 数据库状态机的关键并发控制
一个典型领取任务的逻辑是:
- 找到
READY且available_at <= now的任务; - 通过行锁、乐观锁或原子更新抢占;
- 写入
lease_owner、lease_until、attempt; - 提交事务后执行外部任务;
- 上报结果时校验租约或 fencing token;
- 结果写入只允许作用于当前尝试。
不要只依靠 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 并发与背压
至少分三层限额:
- 租户级:防止一个租户占满平台;
- 工作流级:限制同一流程实例的并行分支;
- 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_definitionworkflow_definition_versionworkflow_instanceworkflow_node_instancenode_attemptworkflow_eventtimerexternal_signaltask_queueworker_leaseartifact/reference9.2 最小字段建议
workflow_definition_version
definition_id、version、schema_version;status:draft/published/retired;definition_json;input_schema、output_schema;created_by、created_at;- 内容 hash,保证版本不可悄悄修改。
workflow_instance
instance_id、definition_id、definition_version;business_key、tenant_id;status、current_revision;input_ref、output_ref;started_at、deadline_at、completed_at;parent_instance_id、correlation_id。
node_attempt
attempt_id、instance_id、node_id;attempt_no、status;queue_name、worker_id、lease_until、fencing_token;idempotency_key;available_at、started_at、finished_at;input_ref、output_ref、error_ref。
workflow_event
event_id、instance_id、sequence;event_type、payload、payload_schema_version;occurred_at、recorded_at、source;dedup_key;- 对
(instance_id, sequence)和必要的dedup_key建唯一约束。
9.3 输入输出存储
小输入、小输出可以存数据库 JSON;大对象、文件和日志应存对象存储,并在工作流中保存引用、hash、大小、媒体类型和保留期限。不要把密码、访问令牌和完整个人敏感数据无差别写入事件历史或日志。
9.4 一致性边界
建议把数据库作为控制面的权威来源:
- 状态推进、租约、事件追加在同一存储中保持原子;
- 队列只承担投递和唤醒;
- worker 结果必须带实例、节点和尝试身份;
- 读模型可以异步最终一致;
- 关键管理操作要记录审计事件。
9.5 分片和归档
初期可按 tenant_id 或 instance_id 做逻辑分区;规模上来后再按时间、租户和实例 hash 组合分片。历史执行数据应设置保留期、冷热分层和归档策略,否则查询索引和存储成本会持续增长。
10. 可观测性、运维与安全
10.1 三类可观测信号
采用 OpenTelemetry Observability Primer中的三大信号思路:
- Metrics:吞吐、成功率、重试率、队列深度、调度延迟、任务耗时、租约过期、定时器延迟;
- Logs:结构化记录实例、节点、尝试、trace、租户和错误码;
- Traces:从 API 请求到工作流实例、任务投递、worker 执行、外部调用的链路。
工作流特有的指标包括:
workflow_started_totalworkflow_completed_totalworkflow_failed_totalworkflow_duration_secondsworkflow_queue_delay_secondsnode_attempt_totalnode_retry_totalnode_timeout_totalorphaned_attempt_totalstuck_instance_total10.2 控制台必须回答的问题
- 这个实例当前在哪个节点?
- 它为什么没有继续?是在等待依赖、队列、定时器还是外部事件?
- 最近一次尝试由哪个 worker 执行?租约何时过期?
- 失败是业务错误、基础设施错误还是未知结果?
- 重试还剩多少次?下一次何时执行?
- 哪些下游服务导致整体变慢?
- 暂停、恢复、取消、重跑、跳过的影响是什么?
10.3 管理操作的安全边界
暂停、取消、重试、跳过节点、修改输入、人工标记成功都属于高风险操作。即使第一版不做角色审批,也至少需要:
- 明确的管理 API;
- 操作者身份和原因;
- 审计记录;
- 防止绕过状态机;
- 对重跑和补偿提供预览;
- 重要操作二次确认或双人机制作为未来扩展。
10.4 多租户安全
如果未来服务多个团队,设计时就引入 tenant_id:
- 所有实例、任务、事件、日志和对象引用带租户边界;
- 查询必须强制带租户过滤;
- worker 队列按租户或资源等级隔离;
- 限制单租户实例数、并发数、历史量和事件大小;
- Secret 通过外部密钥管理系统引用;
- 任务代码不应获得引擎数据库的任意写权限。
11. 参考产品与架构对比
以下对比用于学习架构,不等于直接推荐。请以各项目当前官方文档为准。
11.1 Temporal / Cadence 类持久化执行引擎
- 资料:Temporal Workflows、Activities、Workflow 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 Concepts、Prefect 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 Functions | DSL | 强 | 强 | 中 | 强 | 低 | 托管状态机 |
| Durable Functions | 代码优先 | 强 | 中 | 中 | 强 | 低 | 重放与编排器 |
| Argo | YAML/DSL | 中 | 强 | 强 | 中 | 中 | 容器工作流 |
| Airflow | Python/DAG | 中 | 强 | 中 | 弱~中 | 中 | 数据调度 |
| Dagster | Python/资产 | 中 | 强 | 中 | 中 | 中 | 数据资产 |
| Prefect | Python | 中 | 强 | 中 | 中 | 中 | 动态 flow |
| Conductor | JSON/worker | 强 | 中 | 中 | 强 | 中 | 微服务编排 |
| Camunda | BPMN/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-definitionsPOST /workflow-definitions/{id}/versions/{version}/publishPOST /workflow-instancesGET /workflow-instances/{instanceId}POST /workflow-instances/{instanceId}/signalPOST /workflow-instances/{instanceId}/pausePOST /workflow-instances/{instanceId}/resumePOST /workflow-instances/{instanceId}/cancelGET /workflow-instances/{instanceId}/eventsGET /tasks/pollPOST /tasks/{attemptId}/heartbeatPOST /tasks/{attemptId}/completePOST /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 验收场景
至少实现并验证以下流程:
- 顺序成功:A → B → C;
- 条件分支:根据输入选择 B 或 C;
- 并行汇聚:B/C 并行,全部成功后 D;
- 单任务瞬时失败:退避后成功;
- 永久失败:达到策略后流程失败;
- worker 崩溃:租约过期后重新投递;
- 迟到回报:旧尝试不能覆盖新尝试;
- 重复回报:不会重复推进后续节点;
- 外部事件:流程等待回调后继续;
- 重启恢复:引擎重启后从持久化状态继续;
- 取消:取消后不可再启动新的普通节点;
- 超时:任务超时与工作流总超时行为符合定义。
13.4 最小可用数据表
workflow_definition_versionsworkflow_instancesnode_instancesnode_attemptsworkflow_eventstimersexternal_signals可以先不引入单独的队列表,使用 node_attempts 的 ready_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 周)
建议至少动手:
- 用关系数据库写一个顺序/条件/并行状态机;
- 用队列接入 worker,加入租约和重复投递;
- 加入持久化定时器和外部信号;
- 阅读 Temporal 的 workflow/activity 与事件历史;
- 阅读 Step Functions 的状态语言;
- 阅读 Airflow 或 Dagster 的 DAG/资产模型;
- 在 Argo 中运行一个容器 DAG,观察 Kubernetes 资源边界。
输出:一张“同一流程在不同模型中的表达”对比表,以及一份取舍记录。
15.3 第三阶段:实现可用 MVP(4~8 周)
顺序建议:
- DSL 和静态校验;
- 定义版本和实例创建;
- 状态迁移和事件记录;
- worker 协议和租约;
- 重试、超时、心跳;
- 定时器和外部信号;
- 控制台查询和操作;
- 故障注入与压测。
每加一种能力,都先写状态迁移测试和失败场景测试,不要先堆管理页面。
15.4 第四阶段:生产化(持续迭代)
- 明确 SLO:调度延迟、恢复时间、数据保留和可用性;
- 做容量模型:实例数、节点数、事件数、日志量;
- 做灾备演练和恢复演练;
- 做版本迁移、灰度和回滚;
- 引入多租户、配额、审计和密钥管理;
- 建立工作流定义评审和发布规范。
15.5 调研每个产品时都问同样的问题
- 工作流状态保存在哪里?
- 任务投递是至少一次还是其他语义?
- worker 崩溃后如何恢复?
- 副作用如何保证幂等?
- 长时间等待如何实现?
- 版本发布后,运行中的实例如何处理?
- 事件历史是否完整保存?如何归档?
- 是否支持暂停、取消、重跑和人工修复?
- 调度器如何分片和扩展?
- 数据库、队列和 worker 之间的事实边界是什么?
- 任务代码运行在哪个信任边界?
- 出现未知结果时系统如何避免重复副作用?
16. 结论与决策清单
16.1 对当前项目的推荐结论
- 模型:从数据库驱动的持久化状态机开始,同时保留不可变事件历史;
- 流程表达:先用小型 JSON/YAML DSL,不要一开始实现完整 BPMN;
- 执行方式:worker 拉取 + 租约 + 心跳 + fencing token;
- 可靠性:默认至少一次投递,业务有效一次由幂等键和唯一约束实现;
- 等待:把定时器、外部事件和回调持久化,禁止用线程长期阻塞;
- 数据:控制面数据库是事实来源,队列只负责投递和唤醒;
- 扩展:通过任务类型、worker SDK、子流程和补偿扩展,而不是让引擎直接执行任意业务代码;
- 演进:当流程复杂度和运行时长增长,再引入事件重放、快照、归档或代码优先模型;
- 架构:先模块化单体,验证语义后再按吞吐和故障隔离拆分;
- 验证:把重启、重复投递、迟到回报、未知结果和外部事件乱序列为一等测试场景。
16.2 开始编码前必须写下的 ADR
- 为什么选择状态机而不是完整 BPMN?
- 状态表和事件历史谁是权威?
- 任务投递与业务有效一次如何定义?
- worker 租约和 fencing 如何工作?
- 运行中实例遇到新定义版本怎么办?
- 取消是尽力而为还是强保证?
- 超时后如何处理可能已经成功的外部副作用?
- 大输入输出和敏感数据放在哪里?
- 哪些管理操作允许跳过状态机?
- 什么时候从模块化单体拆成服务?
16.3 最终判断标准
一个工作流引擎是否设计得好,不看它能画出多漂亮的流程图,而看它在以下场景下是否仍然可解释、可恢复、可测试:
- worker 在外部副作用完成后崩溃;
- 消息重复、乱序、迟到或丢失;
- 引擎重启、数据库短暂不可用;
- 一个任务需要等待数天;
- 工作流定义已经升级,但旧实例仍在运行;
- 一个租户突然产生大量并行任务;
- 运维人员需要取消、重试或修复一个异常实例。
如果这些问题能用明确的状态、事件、租约、幂等键和审计记录回答,才说明引擎具备可靠的基础。
17. 参考资料总表
标准与规范
工作流引擎与编排平台
- Temporal Workflows
- Temporal Activities
- Temporal Event History
- Temporal Workflow Definition
- AWS Step Functions — Amazon States Language
- Azure Durable Functions Overview
- Argo Workflows Documentation
- Apache Airflow Core Concepts
- Dagster Concepts
- Prefect Flows
- Conductor OSS Documentation
- Camunda Documentation
分布式系统与可靠性
建议继续检索的主题
- 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。
如果这篇文章对你有帮助,欢迎分享给更多人!
部分信息可能已经过时