Conductor 持久化执行语义(Durable Execution):分布式工作流的每一步状态都可靠落盘
【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductor
本文是 Conductor(一个面向应用与 AI Agent 的事件驱动工作流引擎)的持久化执行(Durable Execution)语义权威指南。文章围绕"每个工作流执行在每个步骤都会被持久化、能够抵御基础设施故障、并保证任务至少一次(at-least-once)投递"这一核心模型展开,完整覆盖持久化内容清单、任务投递保证、故障矩阵、任务状态机、超时与重试配置、工作流级耐久性、重放与恢复以及分布式一致性。读完本文,你将掌握 Conductor 如何在服务器重启、Worker 崩溃、网络分区等各类故障下不丢失执行进度,以及如何据此设计幂等的 Worker 与可安全升级的工作流定义。
什么是 Durable Execution:引擎的可靠性底座
Conductor 本质上是一个面向分布式工作流与持久化 Agent 的执行引擎。所谓 Durable Execution(持久化执行),指的是工作流的每一次执行都在每一步被完整持久化:执行进度不会因进程崩溃、节点重启或机房故障而丢失。结合任务队列、重试与超时机制,引擎能够保证任务至少一次投递,从而让"工作流与 Agent 永不丢失进度"成为可验证的工程承诺。
从代码结构看,这一模型贯穿引擎核心:
- TaskModel.java 定义了任务的数据模型与状态机(
Status枚举),每个任务实例携带scheduledTime、startTime、endTime、updateTime、retryCount、pollCount等时间戳与计数,正是这些字段支撑了后续的持久化、超时判定与重试; - WorkflowModel.java 定义了工作流实例的状态模型;
- DeciderService.java 与 WorkflowSweeper.java 等实现了"decide 求值 + sweeper 清扫"的推进机制。
每一步都持久化:What Persists
当一个工作流执行时,Conductor 会持久化以下四类数据:
- 工作流定义快照(Workflow definition snapshot)——本次执行所使用的定义副本,启动后即不可变(immutable);
- 工作流状态(Workflow state)——状态、输入、输出、correlation ID 以及变量(variables);
- 每一次任务执行(Every task execution)——状态、输入、输出、时间戳、重试次数与 worker ID;
- 任务队列状态(Task queue state)——哪些任务处于已调度、进行中或已完成。
全部状态都会在进入下一步之前写入所配置的持久化存储(Redis、PostgreSQL、MySQL 或 Cassandra)。如果服务器重启,执行将从最后持久化的状态恢复,而不是从头开始。
"定义快照"这一设计意义重大:它意味着运行中的执行与元数据存储(metadata store)解耦。即使工作流定义在运行期间被更新甚至被删除,正在运行的实例依然使用自己内嵌的快照继续执行——这直接支撑了"零停机升级"。
任务投递保证:At-Least-Once Delivery
Conductor 对所有任务提供至少一次投递保证,其循环如下:
- 任务被调度时,放入持久化任务队列(persistent task queue);
- Worker 轮询(poll)并领取任务,任务进入
IN_PROGRESS; - Worker 完成任务后上报
COMPLETED,Conductor 推进工作流; - 若 Worker 失败或崩溃,任务将基于重试与超时配置被重新投递(redelivered)。
一个任务永远不会被静默丢失。如果 Worker 领取了任务却始终不响应,**响应超时(response timeout)**会触发重新投递。
值得注意的是,这同时是"至少一次"而非"恰好一次"——同一任务可能被执行多次。因此,Worker 必须以幂等的方式处理副作用,这一点在文末"这对你的代码意味着什么"一节有进一步说明。
故障矩阵:每种故障场景下引擎的确切行为
下表给出了 Conductor 在各类故障场景下的精确行为:
| 场景 | Conductor 的行为 | 结果 |
|---|---|---|
| Worker 轮询后在开始任何工作前崩溃 | 触发响应超时(response timeout)。任务回到SCHEDULED,由新 Worker 领取。 | 任务自动重试,无数据丢失。 |
| Worker 在产生副作用之后、上报完成之前崩溃 | 触发响应超时。任务被重新投递给另一个 Worker。 | 任务会再次执行。Worker 必须对副作用幂等,或使用任务的updateTime检测重投递。 |
| Worker 上报 FAILED | Conductor 根据重试配置(retryCount、retryDelaySeconds、retryLogic)创建一次新的任务执行。 | 重试至配置上限。重试耗尽后任务进入FAILED,工作流的失败处理逻辑接管。 |
| Worker 上报 FAILED_WITH_TERMINAL_ERROR | 不重试,任务立即终止。 | 工作流失败,或执行配置的failureWorkflow。 |
| 工作流执行期间服务器重启 | 重启后,sweeper 服务从持久化存储拾取进行中的工作流并重新求值(re-evaluate)。 | 从最后持久化的状态恢复执行,无需人工干预。 |
| 跨多次部署的长时间等待 | WAIT 与 HUMAN 任务在持久化存储中保持IN_PROGRESS。计时器或信号的解析是持久的。 | 当等待时长耗尽或信号到达时(即使是在多次部署之后的数天),任务完成,工作流继续推进。 |
| 暂停中的工作流收到信号/Webhook | Task Update API 或事件处理器将 WAIT/HUMAN 任务置为COMPLETED并提供输出。 | 工作流立即恢复,信号载荷可作为任务输出使用。 |
| 运行期间工作流定义被更新 | 运行中的执行继续使用启动时拍摄的定义快照,新执行使用新定义。 | 定义变更不影响任何运行中的执行,实现零停机升级。 |
| 运行期间工作流版本被删除 | 运行中的执行与元数据存储解耦,继续使用内嵌的定义快照。 | 现有执行正常完成,只有新启动的实例受影响。 |
| Worker 与服务器之间网络分区 | Worker 的更新无法到达服务器,触发响应超时,任务重新入队。 | 分区恢复后,新 Worker(或同一个 Worker)重新领取任务。 |
这张矩阵是本文档的灵魂,它用一张表穷举了"任何环节出问题都不会丢进度"的承诺边界:要么自动重试,要么转交失败处理,要么等待恢复——但绝不会出现"任务凭空消失"的状态。
任务状态机:Task State Transitions
每个任务遵循以下状态机:
SCHEDULED ──→ IN_PROGRESS ──→ COMPLETED │ │ │ ├──→ FAILED ──→ SCHEDULED (retry) │ │ │ ├──→ FAILED_WITH_TERMINAL_ERROR │ │ │ └──→ TIMED_OUT ──→ SCHEDULED (retry) │ └──→ CANCELED (workflow terminated)**终态(Terminal states)**包括:COMPLETED、FAILED(重试耗尽后)、FAILED_WITH_TERMINAL_ERROR、CANCELED、COMPLETED_WITH_ERRORS(可选任务 optional tasks)。
每一次状态转移都会在任何后续动作发生之前被持久化。
在源码中,这套状态机被建模为 TaskModel.Status 枚举,并且每个状态显式标注了三个关键属性:
| 状态 | terminal | successful | retriable |
|---|---|---|---|
IN_PROGRESS | 否 | 是 | 是 |
CANCELED | 是 | 否 | 否 |
FAILED | 是 | 否 | 是 |
FAILED_WITH_TERMINAL_ERROR | 是 | 否 | 否 |
COMPLETED | 是 | 是 | 是 |
COMPLETED_WITH_ERRORS | 是 | 是 | 是 |
SCHEDULED | 否 | 是 | 是 |
TIMED_OUT | 是 | 否 | 是 |
SKIPPED | 是 | 是 | 否 |
这一枚举定义直接印证了文档中的行为描述:FAILED_WITH_TERMINAL_ERROR的retriable=false,因此引擎不会对它重试;FAILED与TIMED_OUT的retriable=true,因此会按配置重试后重新回到SCHEDULED。isRetriable()、isSuccessful()、isTerminal()这三个方法就是引擎在 decide 流程中判断"是否该重试、是否算成功、是否已到终态"的依据。
超时与重试配置:每任务级参数
耐久性可以通过任务定义(task definition)按任务单独配置,详见 taskdef.md。核心参数如下:
| 参数 | 作用 |
|---|---|
timeoutSeconds | 任务到达终态所允许的最大墙钟时间(wall-clock time)。 |
responseTimeoutSeconds | 在重新入队前,等待 Worker 状态更新的最大时间。 |
pollTimeoutSeconds | 一个已调度任务在被轮询前等待的最大时间,超时即触发超时。 |
retryCount | 失败或超时时的重试次数。 |
retryLogic | FIXED、EXPONENTIAL_BACKOFF或LINEAR_BACKOFF。 |
retryDelaySeconds | 重试之间的基础延迟。 |
timeoutPolicy | RETRY、TIME_OUT_WF或ALERT_ONLY。 |
从源码 TaskDef.java 可以看到这些枚举与默认值的真实定义:
public enum TimeoutPolicy { RETRY, TIME_OUT_WF, ALERT_ONLY } public enum RetryLogic { FIXED, EXPONENTIAL_BACKOFF, LINEAR_BACKOFF }其默认值分别为:
retryCount默认为3;timeoutPolicy默认为TIME_OUT_WF(即任务超时后直接判定工作流超时失败);retryLogic默认为FIXED;retryDelaySeconds默认为60秒;timeoutSeconds无默认值,需显式配置(且带@NotNull校验约束)。
理解这几个默认值有助于避免"我以为不会重试、结果重试了 3 次"或"我以为会重试、结果工作流直接超时失败"之类的配置误区。其中retryDelaySeconds是三种重试逻辑共用的"基础延迟":FIXED每次固定等待该时长;EXPONENTIAL_BACKOFF按指数递增;LINEAR_BACKOFF按线性递增。responseTimeoutSeconds与pollTimeoutSeconds则共同决定了"多久判定一个 Worker 失联、多久判定一个任务无人领取",是故障矩阵中"Worker 崩溃后自动重投递"得以实现的计时器基础。
工作流级耐久性:超越单任务
除单个任务外,Conductor 还提供工作流级别的耐久能力:
- 补偿流(Compensation flows):配置一个
failureWorkflow,当主工作流失败时自动运行,并携带完整上下文(失败原因、失败任务 ID、工作流执行数据); - 暂停与恢复(Pause and resume):任意运行中的工作流可通过 API 暂停并在之后恢复,状态被完整保留;
- 重启、重跑与重试(Restart / rerun / retry):详见下文"重放与恢复"一节;
- 版本化(Versioning):多个工作流版本可以并发运行;运行中的执行对定义变更不可变;重启时可选地使用最新定义。
这里的failureWorkflow是实现Saga 补偿模式的官方入口:主流程失败后,自动触发补偿流程撤销已完成的副作用,而补偿流程本身同样享受整套持久化保证。
重放与恢复:Replay and Recovery
每一个工作流执行都是完全可重放的(fully replayable)。Conductor 保留了完整的执行图——每个任务的输入、输出与状态——因此你可以随时重新执行工作流。
| 操作 | 作用 | 适用场景 |
|---|---|---|
| Restart(重启) | 从开头重新执行整个工作流 | 定义已变更,需要一次干净的执行 |
| Rerun(重跑) | 从某个特定任务开始重新执行,复用之前任务的输出 | 修复中间某个任务,而无需重跑全部 |
| Retry(重试) | 重试最后一个失败的任务,并从该点继续 | 瞬时故障、外部依赖当时不可用 |
这三个操作都可以作用于任意终态(COMPLETED、FAILED、TIMED_OUT、TERMINATED)的工作流,并且可以无限期使用——因为 Conductor 完整保留了执行图。Restart 还可以选择性地使用最新的工作流定义,这样你可以在修复定义中的 bug 后立刻重放执行。
分布式一致性:多节点部署下的正确性
在多节点部署中,Conductor 通过以下机制保证一致性:
- 分布式锁(Distributed locking):在整个集群中,每个工作流同一时刻只有一个
decide求值在运行(可插拔实现:Zookeeper、Redis); - 栅栏令牌(Fencing tokens):防止持有过期锁的节点提交过期更新;
- 持久化队列(Persistent queues):任务队列在节点故障后依然存活。支持可配置的分片策略(round-robin 或 local-only),在分布性与一致性之间做权衡。
分布式锁配置详见部署指南中的 locking 一节。在源码层面,这对应 WorkflowReconciler.java、WorkflowSweeper.java 与 ExecutionLockService.java 等组件:sweeper 定期从存储中拾取未决工作流,在分布式锁的保护下执行 decide,确保同一工作流的推进在任何时刻只发生在一个节点上,从而避免"双份推进导致的状态错乱"。
这对你的代码意味着什么
- Worker 应该是幂等的。由于至少一次投递保证,任务可能被执行不止一次。请将 Worker 设计为能够安全处理重投递。
- 你不需要自己构建重试逻辑。Conductor 负责重试、超时与重新入队。你的 Worker 只需上报成功或失败。
- 长时间运行的流程是安全的。使用 WAIT 与 HUMAN 任务来处理跨越数分钟到数天的暂停,状态在多次部署之间保持持久。
- 定义变更是安全的。可以随时更新工作流定义而不影响正在运行的执行。以零停机的方式逐步发布新版本。
作为补充,对于"Worker 在副作用之后崩溃"的场景,文档给出了两条工程路径:一是让 Worker 对副作用幂等(重复执行同一副作用的结果相同);二是利用任务上的updateTime字段(见 TaskModel.java 中的字段定义)检测"这个任务是否已经处理过一次",从而在重投递时跳过已完成的副作用。这一字段随任务状态一同持久化,正是 Durable Execution 模型中"可检测的重投递"与"业务幂等"之间的衔接点。
总而言之,Conductor 的持久化执行语义可以浓缩为一句话:每一步都落盘,任何故障都有明确的、可预期的行为路径。基于这一模型,你可以放心地把跨机器、跨进程、跨部署的工作流编排任务交给引擎,而把精力集中在业务逻辑本身。
【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductor
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考