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
导读
本文是 Conductor 中把"本地跑通的工作流"升级为"可长期运行的生产服务"的实操路线图。它围绕 docs/devguide/workflows/production-path.md 给出的五个生产阶段——定义契约、明确故障策略、测试真实边界、安全部署、运维执行——展开,并深入结合本仓库的源码与配套文档,说明每一条建议背后的实现机制。读完本文,你将掌握如何为工作流建立可被调用方依赖的输入/输出契约、如何为每个外部副作用设计可重复的故障行为、如何在发布前测试真实边界、如何安全地滚动版本,以及如何在事故发生时定位、检查并恢复一次执行。
前置条件:建议先完成 docs/quickstart/first-workflow.md 中基于内置系统任务(HTTP、JSON_JQ_TRANSFORM)的两步工作流,理解"定义注册—启动执行—查看结果"的基本闭环,再进入本文的生产化改造。
生产路径全景:五阶段路线图
原文档将生产化过程归纳为五个阶段,每一阶段都对应一个可验证的产出:
| 阶段 | 核心动作 | 交付物 |
|---|---|---|
| 1. 定义契约 | 把工作流定义与outputParameters当作 API 对待 | 调用方可知的输入/输出契约 |
| 2. 明确故障策略 | 为每个外部副作用设计重试、幂等与超时行为 | 有界且可预期的失败语义 |
| 3. 测试真实边界 | 用代表性输入测试已注册的定义,而非只测 worker 函数 | 覆盖成功/重试/终态失败/超时/幂等的验证记录 |
| 4. 安全部署 | 先部署 worker 与任务定义,再切流量;按版本滚动 | 无中断的版本发布与排空 |
| 5. 运维执行 | 约定命名、版本、关联 ID 与负责人,监控并演练恢复 | 可被定位、检查与恢复的执行体系 |
这与 docs/devguide/workflows/index.md 中描述的工作流生命周期一一对应:Build(构建定义)→ Register(注册版本)→ Trigger(启动触发)→ Execute(持久化执行)→ Observe(检查运维)→ Evolve(版本演进)。生产路径本质上就是让这一生命周期中的每个环节都显式化、可审计、可恢复。
关键前提是 Conductor 的定义与执行分离模型:工作流定义声明了任务顺序与数据流,而工作流执行是蓝图的单次运行,拥有自己的 ID、输入和历史。因为二者分离,修改定义永远不会改写已运行执行的历史——这正是版本演进与安全部署的基础,具体语义见 docs/devguide/concepts/workflows.md 与 docs/architecture/durable-execution.md。
阶段一:把工作流定义当作 API 来定义契约
契约的两个锚点:输入声明与 outputParameters
在生产中,工作流定义的第一职责是成为一份稳定的接口契约。一个最小但完整的定义应显式声明:
inputParameters:工作流期望的输入键清单;outputParameters:输出键到表达式(通常引用任务输出)的映射,调用方据此消费结果;tasks:任务配置序列,任务之间通过taskReferenceName引用彼此的输出。
参考 docs/devguide/how-tos/Workflows/creating-workflows.md 中的最小示例结构:
{ "name": "order_flow", "version": 1, "schemaVersion": 2, "inputParameters": ["orderId"], "tasks": [ { "name": "process_order", "taskReferenceName": "process_order_ref", "type": "SIMPLE", "inputParameters": { "orderId": "${workflow.input.orderId}" } } ], "outputParameters": { "status": "${process_order_ref.output.status}" } }原文档强调的要点是:把工作流定义连同其outputParameters视为一个 API。这意味着三件事:
- 文档化必需输入:在定义中用
inputParameters声明期望的输入键,并配套文档说明格式与约束。 - 在边界处校验:对非法请求要能拒绝,而不是让非法输入流经整条任务链后在某个深处以晦涩的方式失败。
- 保持输出稳定:调用方依赖
outputParameters的键名与结构,跨版本保持稳定;当变更不向后兼容时,注册新版本而不是原地修改活跃定义。
动态引用语法:契约的数据流基础
输入输出契约在运行时通过动态引用表达式完成数据传递,语法为${type.jsonpath}。参考 docs/devguide/how-tos/Tasks/task-inputs.md,核心引用类型包括:
| 引用形式 | 含义 |
|---|---|
${workflow.input.<key>} | 当前工作流的输入参数 |
${workflow.variables.<key>} | 工作流变量(由 SET_VARIABLE 任务写入) |
${<taskReferenceName>.output.<key>} | 某任务(按引用名)的输出参数 |
${<taskReferenceName>.input.<key>} | 某任务的输入参数 |
${workflow.workflowId}、${workflow.correlationId}、${workflow.version} | 执行 ID、关联 ID、版本号 |
在定义契约阶段,务必逐个核对outputParameters中每个表达式引用的任务输出是否确实存在、类型是否一致。若引用表达式格式错误,或引用的任务输出在引用点尚未解析,参数值会变成错误数据或 null——这类问题应在发布前通过真实执行暴露,而不是留给调用方在生产中踩坑。
向后不兼容时的新版本注册
当输入、输出、任务顺序或失败语义发生调用方可观察的变化时,不要原地修改正在被调用的定义版本,而是递增version字段并注册新版本。版本化的运行时语义是:每次执行在启动时引用工作流定义的一个快照,此后对定义的任何修改都不会影响已运行执行。这一机制在 docs/devguide/how-tos/Workflows/versioning-workflows.md 中有完整说明,并可从 docs/architecture/durable-execution.md 的故障矩阵中得到印证:"运行中的执行继续使用启动时的定义快照,新执行使用更新后的定义,零停机升级"。
阶段二:把故障策略显式化
原文档给出的原则是:对每一个外部副作用,都要先回答两个问题——它是否安全重试?如何让它幂等?然后再决定重试行为与超时配置,并在需要业务回滚时引入故障工作流(failure workflow)或补偿逻辑。
至少一次交付与幂等 Worker
Conductor 对任务采用at-least-once(至少一次)交付:网络分区、worker 重启、响应超时都可能导致同一任务被投递多次。因此 worker 必须以幂等为默认设计。常见模式(见 docs/devguide/bestpractices.md):
| 模式 | 适用场景 |
|---|---|
| 幂等键 | 将workflowId + taskId作为唯一键传给下游服务,由下游去重 |
| Upsert 而非 Insert | 重复写入收敛到同一状态 |
| Check-then-act | 执行前查询当前状态,已完成则跳过 |
| 幂等 HTTP 方法 | 下游支持时优先 PUT 而非 POST |
workflowId与taskId的组合在每次任务执行尝试内唯一,是理想的幂等键来源。幂等不是可选项:只要 worker 可能被重复调用,就必须保证"执行两次与执行一次结果相同且无副作用"。
超时与重试参数的刻意配置
任务定义中的超时参数是故障边界的主要控制面,完整字段契约见 docs/documentation/configuration/taskdef.md。核心规则与推荐起点:
responseTimeoutSeconds<timeoutSeconds:前者是心跳窗口,用于探测无响应的 worker;后者是任务的整体 SLA。timeoutSeconds为 0 表示不设超时;responseTimeoutSeconds为 0 表示禁用响应超时(仅适合由外部完成的任务,如 WAIT/HUMAN)。retryLogic支持FIXED、EXPONENTIAL_BACKOFF、LINEAR_BACKOFF,配合retryDelaySeconds计算基础延迟,maxRetryDelaySeconds封顶延迟,backoffJitterMs引入抖动防止惊群。
参考 docs/devguide/bestpractices.md 的推荐配置矩阵:
| 任务模式 | responseTimeoutSeconds | timeoutSeconds | timeoutPolicy | retryCount |
|---|---|---|---|---|
| API 调用(预期 < 5s) | 10 | 30 | RETRY | 3 |
| ML 推理 | 120 | 300 | RETRY | 1 |
| 人工审批 | 0(禁用) | 86400 | ALERT_ONLY | 0 |
| 批处理 | 600 | 3600 | TIME_OUT_WF | 0 |
| 快速数据转换 | 5 | 15 | RETRY | 3 |
超时策略决定超时后的动作:
RETRY:按retryCount重试任务(瞬态失败如网络、外部 API);TIME_OUT_WF:立即失败整个工作流(任务至关重要且重试无意义,如过期批窗口);ALERT_ONLY:标记任务超时但保持工作流运行(人工参与或外部完成信号的任务)。
重试逻辑决定两次重试间的延迟:FIXED恒定间隔;EXPONENTIAL_BACKOFF按retryDelaySeconds × 2^attemptNumber递增,适合限流 API 与过载服务;LINEAR_BACKOFF按retryDelaySeconds × attemptNumber线性递增,适合中等恢复场景。最终延迟公式为delay = clamp(computedDelay, 0, maxRetryDelaySeconds) + random(0, backoffJitterMs) ms,见 docs/documentation/configuration/taskdef.md 的重试逻辑小节。
cookbook/task-timeouts-and-retries.md 提供了几个可直接注册的完整任务定义配方,例如带封顶与抖动的指数退避、长任务心跳租约续期(callbackAfterSeconds)、以及用totalTimeoutSeconds强制硬 SLA——它是独立于retryCount的墙钟总预算,谁先到达谁生效。当业务操作有最大可接受时长时,还应给工作流本身设置timeoutSeconds与timeoutPolicy(工作流级支持TIME_OUT_WF与ALERT_ONLY),把整个执行也约束在有界时间内。
失败工作流与 Saga 补偿
当后续失败需要业务回滚时,用failureWorkflow声明补偿工作流:主工作流失败后,Conductor 自动触发它。可在主定义中同时指定failureWorkflowVersion以固定补偿流的版本。从源码看,这一路径由核心执行引擎直接支持——WorkflowExecutor.java 中terminateWorkflow(...)的重载签名携带failureWorkflow与failureWorkflowVersion参数,负责在终止失败工作流时联动触发补偿流。
默认情况下,以下参数会作为输入传给失败工作流(见 docs/devguide/how-tos/Workflows/handling-errors.md):
| 参数 | 内容 |
|---|---|
reason | 工作流失败原因 |
workflowId | 失败执行的 ID |
failureStatus | 失败工作流的状态 |
failureTaskId | 失败任务的执行 ID |
failedWorkflow | 失败工作流的完整执行 JSON |
补偿任务应标记optional: true(针对失败前可能未执行到的步骤),并利用failedWorkflow.tasks[...].output中的事务 ID、资源句柄等元数据完成反向回滚——先取消发货、再恢复库存、最后退款,逆序撤销已完成步骤。这正是 Saga 模式在 Conductor 中的落地形态,完整可运行示例见 docs/devguide/how-tos/Workflows/handling-errors.md 中的order_processing/order_compensation双工作流示例。
验证故障设计:强制一次瞬态失败
原文档给出了明确的验收手段:强制一次瞬态任务失败,确认预期的重试、超时或补偿路径在执行中可见。例如让一个SIMPLEworker 故意抛错,观察任务按retryLogic重试、最终耗尽retryCount后转入FAILED,随后failureWorkflow被触发并携带reason、failedWorkflow等上下文。对不可重试的业务错误,worker 应返回FAILED_WITH_TERMINAL_ERROR状态跳过全部重试立即终态——例如余额不足的支付,重试永远不会成功。
阶段三:测试真实边界
原文档强调:测试已注册的定义(用代表性输入),而不只是隔离测试 worker 函数。原因是定义级别的编排逻辑、数据流、分支决策只有端到端执行才能验证。测试应覆盖:
- 成功路径;
- 可重试的瞬态失败;
- 终态业务失败;
- 超时;
- 每个副作用的幂等行为。
三层测试法
docs/devguide/how-tos/Workflows/testing-workflows.md 给出了三层递进的测试结构:
- Schema 校验:
POST /api/metadata/workflow/validate,成功返回空200 OK。它只校验元数据与图规则,不证明 worker 在轮询、外部端点可达。 - Mock 化编排测试:
POST /api/workflow/test用按引用名提供的 mock 输出来模拟任务,验证分支与数据流决策,不调用真实 worker。每个引用名对应一个列表,因为循环或重试可能消耗多个 mock;mock 上的executionTime、queueWaitTime可模拟超时行为。 - 真实边界执行:注册定义后启动一次真实执行,验证 worker、队列、持久化与并发行为。使用真实依赖或 Testcontainers 时效果最佳。
用测试关联 ID 做端到端验收
原文档给出的验收步骤值得固化到发布清单中:
- 用测试关联 ID(correlation ID)启动工作流;
- 检查完整执行(
conductor workflow get-execution <workflow-id> -c); - 断言输出契约与终态状态。
启动时可通过 CLI 携带关联 ID 与固定版本(见 docs/devguide/how-tos/Workflows/starting-workflows.md):
conductor workflow create workflow.json conductor workflow start -w order_flow --version 1 \ --correlation order-test-001 -i '{"orderId":"order-001"}' conductor workflow get-execution <workflow-id> -c关联 ID 便于在 docs/devguide/how-tos/Workflows/searching-workflows.md 中按业务维度检索执行,但它不保证唯一,也不能替代 workflow ID 作为执行的主键。需要特别注意:SIMPLE任务若缺少已注册的任务定义或轮询中的 worker,执行会一直停留在队列中——mock 测试无法发现这类部署缺口,必须由真实执行暴露。
阶段四:安全部署定义与 Worker
顺序:先 worker 与任务定义,后生产流量
原文档的部署铁律:在把生产流量路由到某个工作流之前,先部署好它依赖的 worker 代码与任务定义。因为 Conductor 采用至少一次交付,worker 必须幂等;因为SIMPLE任务需要"已注册的任务定义 + 轮询该精确任务类型的 worker"两者齐备,缺一则任务停留在队列、工作流不推进(docs/devguide/concepts/workflows.md)。部署前用conductor taskDef list核对每个SIMPLE任务的依赖是否齐备(docs/devguide/how-tos/Workflows/creating-workflows.md)。
按版本滚动发布
新版本工作流的发布流程(参考 docs/devguide/how-tos/Workflows/versioning-workflows.md):
- 递增
version并注册新版本(conductor workflow create workflow.json),而非覆盖生产调用方使用的版本; - 校验与 mock 测试后,用固定版本启动一次金丝雀执行:
conductor workflow start -w <name> --version 2 -i '{"orderId": "test-1"}'; - 刻意地把调用方、调度器与父工作流引用迁移到新版本;
- 对比两个版本的完成率、失败率、延迟与输出;
- 保留旧版本直至调用方完成迁移、旧执行的排空不再需要重启/重放支持。
由于定义变更不影响运行中的执行,蓝绿发布天然成立:新执行走版本 N+1,存量版本 N 的执行继续在自己的定义快照上排空至完成。若版本 N+1 出问题,可停止向它启动新执行、让存量失败或终止、回到从未被修改过的版本 N 继续服务。因为 worker 与工作流定义解耦,可以独立于 worker 部署回滚工作流版本。
运行中执行升级的特殊情况
定义变更从不影响运行中执行,因此需要显式升级时,方式是"终止 + 以最新定义重启"。UI 上依次选择Actions > Terminate、Actions > Restart with latest definitions;API 则用批量终止 + 批量重启:
curl -X POST 'http://localhost:8080/api/workflow/bulk/terminate' \ -H 'Content-Type: application/json' \ -d '["<workflow-id-1>", "<workflow-id-2>"]' curl -X POST 'http://localhost:8080/api/workflow/bulk/restart?useLatestDefinitions=true' \ -H 'Content-Type: application/json' \ -d '["<workflow-id-1>", "<workflow-id-2>"]'不传useLatestDefinitions=true时,重启使用各自原始的定�义快照,不会发生升级。仓库中 WorkflowBulkServiceImpl.java 等实现为批量重启提供了useLatestDefinitions语义支持。必须注意:终止并重启会重复副作用——除非工作流幂等或已定义补偿,否则优先让运行中执行在自己的快照上自然完成。
阶段五:运维执行
建立执行的可发现性基线
原文档要求为运维者准备四件套:工作流名、版本、关联 ID 约定、负责人。在此基础上补充两点:
- 启动执行时尽量固定版本(省略 version 会默认取服务器最新注册版本,牺牲可预测性);
- 为执行关联业务意义的 correlation ID,作为跨系统的追踪线索。
监控什么
运维核心是"队列—任务—工作流"三层的可观测性(docs/devguide/how-tos/Workers/scaling-workers.md):
| 指标 | 含义 | 告警参考 |
|---|---|---|
| 任务队列深度 | 未处理任务积压 | 持续 5 分钟增长 |
| 任务轮询数(按类型) | worker 是否在活跃轮询 | 降为零 |
| 工作流失败率 | 以 FAILED 终止的执行占比 | 15 分钟窗口内 > 5% |
| 任务响应时间 p99 | 与响应超时的接近程度 | >responseTimeoutSeconds的 80% |
| worker 线程利用率 | worker 是否饱和 | 持续 10 分钟 > 90% |
| 外部负载存储错误 | S3/GCS 写入失败 | 任何非零计数 |
队列监控可通过 UI(Home > Task Queues)、CLI(conductor task list/conductor task get <TASK_NAME>)或 API(/tasks/queue/sizes、/tasks/queue/polldata)获取;Prometheus 指标如task_queue_depth、task_completed_seconds_count、task_queue_wait_time_seconds均带taskType标签,可支撑按任务类型的告警与自动扩缩容。
事故恢复流程:先检查、后恢复
事故处置的原则是:先检查失败任务,再决定重试。恢复动作包括 pause/resume(暂停/恢复)、rerun(从指定任务重跑)、restart(从头重启)与 terminate(终止),按业务策略选择。参考 docs/devguide/how-tos/Workflows/debugging-workflows.md:
- 从持久化执行开始:
conductor workflow get-execution <workflow-id> -c; - 定位
FAILED/TIMED_OUT/ 终态任务,记录其reasonForIncompletion、输入、输出、worker ID 与重试次数; - 先修复底层的 worker、依赖、凭据或定义问题,再改变执行状态;
- 只在确定安全时才重试(重试/重跑可能重复副作用)。
reasonForIncompletion是诊断的关键字段:worker 返回FAILED/FAILED_WITH_TERMINAL_ERROR时由 worker 写入;HTTP 任务在非 2xx 时记录响应体;引擎则在超时、终止与定义错误时生成模板化消息(例如Task timed out after {elapsed} seconds. Timeout configured as {timeoutSeconds} seconds...)。恢复选项对比:
| 动作 | 行为 | 适用场景 |
|---|---|---|
| Restart with Current Definitions | 用原始定义从头重跑 | 定义已变更但想按原定义执行 |
| Restart with Latest Definitions | 用最新定义从头重跑 | 定义已修复,需按最新定义执行 |
| Rerun from a specific task | 从指定任务重跑,复用此前任务输出 | 中间任务失败,不想重跑前置全部步骤 |
| Retry - From failed task | 从最后失败任务重试 | 瞬态失败、外部依赖临时不可用 |
CLI 对应命令:conductor workflow retry <workflow-id>、conductor workflow restart <workflow-id>、conductor workflow rerun <workflow-id> --task-id <task-id>。恢复后运行conductor workflow status <workflow-id>验证预期任务在运行或到达预期终态。restart/rerun/retry 三种操作对任意终态执行(COMPLETED、FAILED、TIMED_OUT、TERMINATED)都可用且长期可用——Conductor 保留完整执行图,数月后仍可重放任意执行。
恢复演练
原文档给出一个可直接执行的演练脚本:
- 故意让一个工作流保持等待,或使某个可重试任务失败;
- 通过 searching-workflows(如
conductor workflow search -w order_processing -s FAILED -c 20)找到它; - 用 debugging-workflows 的恢复控制将其恢复。
演练的价值在于:让运维者熟悉"搜索定位 → 检查失败上下文 → 选择恢复动作 → 验证终态"这条肌肉记忆,确保真实事故时能在有界时间内完成处置。
下一步:平台部署与 AI 工作流
原文档在生产路径之后给出了两条延伸线:
- 平台级部署与存储选型:继续阅读 docs/devguide/running/deploy.md(API Server、Decider、Sweeper、事件处理器、数据库/队列/索引/分布式锁的架构角色,以及 Docker Compose 各后端组合)与 docs/architecture/durable-execution.md(持久化内容、故障矩阵、任务状态机、分布式一致性)。
- AI 工作流 / Agent 基线增强:若你的工作流面向 AI 或 Agent,在本文的生产基线上叠加 docs/devguide/ai/production-agent-architecture.md 的控制面——父工作流持有业务过程、Agent 在显式执行边界之后运行、不可逆操作必须经父级校验与审批;用
failureWorkflow做副作用补偿、用DO_WHILE循环条件做迭代/预算上限、用HUMAN任务做持久化审批门禁,并把关联 ID 与幂等键带入外部效果。
无论是普通业务工作流还是 Agent 工作流,本文的五阶段路径都适用:契约先行、故障有界、测试真实、部署安全、运维可恢复——这正是 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),仅供参考