尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Conductor 持久化执行语义(Durable Execution):分布式工作流的每一步状态都可靠落盘

Conductor 持久化执行语义(Durable Execution):分布式工作流的每一步状态都可靠落盘 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 DeliveryConductor 对所有任务提供至少一次投递保证其循环如下任务被调度时放入持久化任务队列persistent task queueWorker 轮询poll并领取任务任务进入IN_PROGRESSWorker 完成任务后上报COMPLETEDConductor 推进工作流若 Worker 失败或崩溃任务将基于重试与超时配置被重新投递redelivered。一个任务永远不会被静默丢失。如果 Worker 领取了任务却始终不响应**响应超时response timeout**会触发重新投递。值得注意的是这同时是至少一次而非恰好一次——同一任务可能被执行多次。因此Worker 必须以幂等的方式处理副作用这一点在文末这对你的代码意味着什么一节有进一步说明。故障矩阵每种故障场景下引擎的确切行为下表给出了 Conductor 在各类故障场景下的精确行为场景Conductor 的行为结果Worker 轮询后在开始任何工作前崩溃触发响应超时response timeout。任务回到SCHEDULED由新 Worker 领取。任务自动重试无数据丢失。Worker 在产生副作用之后、上报完成之前崩溃触发响应超时。任务被重新投递给另一个 Worker。任务会再次执行。Worker 必须对副作用幂等或使用任务的updateTime检测重投递。Worker 上报 FAILEDConductor 根据重试配置retryCount、retryDelaySeconds、retryLogic创建一次新的任务执行。重试至配置上限。重试耗尽后任务进入FAILED工作流的失败处理逻辑接管。Worker 上报 FAILED_WITH_TERMINAL_ERROR不重试任务立即终止。工作流失败或执行配置的failureWorkflow。工作流执行期间服务器重启重启后sweeper 服务从持久化存储拾取进行中的工作流并重新求值re-evaluate。从最后持久化的状态恢复执行无需人工干预。跨多次部署的长时间等待WAIT 与 HUMAN 任务在持久化存储中保持IN_PROGRESS。计时器或信号的解析是持久的。当等待时长耗尽或信号到达时即使是在多次部署之后的数天任务完成工作流继续推进。暂停中的工作流收到信号/WebhookTask 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 枚举并且每个状态显式标注了三个关键属性状态terminalsuccessfulretriableIN_PROGRESS否是是CANCELED是否否FAILED是否是FAILED_WITH_TERMINAL_ERROR是否否COMPLETED是是是COMPLETED_WITH_ERRORS是是是SCHEDULED否是是TIMED_OUT是否是SKIPPED是是否这一枚举定义直接印证了文档中的行为描述FAILED_WITH_TERMINAL_ERROR的retriablefalse因此引擎不会对它重试FAILED与TIMED_OUT的retriabletrue因此会按配置重试后重新回到SCHEDULED。isRetriable()、isSuccessful()、isTerminal()这三个方法就是引擎在 decide 流程中判断是否该重试、是否算成功、是否已到终态的依据。超时与重试配置每任务级参数耐久性可以通过任务定义task definition按任务单独配置详见 taskdef.md。核心参数如下参数作用timeoutSeconds任务到达终态所允许的最大墙钟时间wall-clock time。responseTimeoutSeconds在重新入队前等待 Worker 状态更新的最大时间。pollTimeoutSeconds一个已调度任务在被轮询前等待的最大时间超时即触发超时。retryCount失败或超时时的重试次数。retryLogicFIXED、EXPONENTIAL_BACKOFF或LINEAR_BACKOFF。retryDelaySeconds重试之间的基础延迟。timeoutPolicyRETRY、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默认为3timeoutPolicy默认为TIME_OUT_WF即任务超时后直接判定工作流超时失败retryLogic默认为FIXEDretryDelaySeconds默认为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),仅供参考
返回列表