AI 工程基础体系 · 第 23/100 篇。内容覆盖机器学习、深度学习与生成式 AI;模型、数据、评测、权限和成本会作为同一生产系统处理。

AI 工作流与多 Agent 编排:DAG、事件、并发、重试和一致性

在一个生产级 AI 系统中,“让模型调用几个工具”只是局部能力。真正困难的是:多个模型、工具、数据集、评测任务和人工审批如何组成可恢复的流程;一个步骤失败后哪些步骤需要重做;两个 Agent 同时修改同一份状态时谁获胜;外部副作用已经发生但本地进程崩溃后如何避免重复执行。

因此,AI 工作流应被视为一种分布式系统,而不是一段更长的 Prompt。机器学习训练、深度学习推理、生成式 AI Agent、RAG 检索、评测和发布流程都可以使用同一套抽象:

输入状态任务事件状态变更\text{输入} \rightarrow \text{状态} \rightarrow \text{任务} \rightarrow \text{事件} \rightarrow \text{状态变更}

其中模型只是某些任务的执行器。模型输出不稳定、工具有副作用、网络会重试、消息可能重复、权限会变化、成本需要预算,这些因素共同决定了工作流的可靠性。

一、先区分四个概念:任务、工作流、Agent 与编排器

**任务(task)**是一次具有明确输入和输出的执行单元。例如:

  • 从文档中抽取结构化字段;
  • 调用向量数据库检索候选文档;
  • 使用分类模型预测风险等级;
  • 让 LLM 生成答案草稿;
  • 运行离线评测;
  • 请求人工批准发布模型。

任务可以由普通代码、机器学习模型、深度学习服务、LLM 或外部工具完成。任务不一定是 Agent。

Agent通常是具有以下能力的运行单元:

  1. 接收目标和上下文;
  2. 根据当前状态决定下一步;
  3. 调用工具或其他 Agent;
  4. 读取结果并继续推理;
  5. 在满足终止条件、失败条件或人工确认条件时结束。

Agent 的核心不是“使用了 LLM”,而是它拥有一定的决策权。一个固定的 model -> parser -> database 流程是工作流;一个根据工具返回结果决定是否继续检索、改写查询或转交专门 Agent 的组件,才更接近 Agent。

**工作流(workflow)**是任务、状态、依赖、事件和终止条件的整体定义。它回答“系统应该经过哪些步骤”。

**编排器(orchestrator)**负责把工作流变成实际执行:

  • 持久化运行状态;
  • 判断哪些任务已经满足依赖;
  • 分配任务给 Worker 或 Agent;
  • 控制并发;
  • 处理超时、重试和取消;
  • 记录事件和观测数据;
  • 在人工审批或权限检查处暂停;
  • 在恢复后继续执行。

因此,Agent 可以是工作流中的一个节点;多个 Agent 也可以由一个编排器协调。不要把“多 Agent”误解成“多个模型并发调用”——并发调用只是执行策略,编排还包括依赖、状态和故障语义。

二、DAG:把“谁依赖谁”形式化

1. DAG 的定义

DAG 是有向无环图(Directed Acyclic Graph)。将工作流表示为:

G=(V,E)G=(V,E)

其中:

  • VV 是任务节点集合;
  • EE 是有向边集合;
  • (u,v)E(u,v)\in E 表示任务 vv 依赖任务 uu 的输出。

如果存在一条路径从节点回到自身,则图存在环,不是 DAG。

例如一个问答系统可以表示为:

flowchart LR
    A[接收问题] --> B[查询改写]
    A --> C[权限检查]
    B --> D[检索]
    C --> D
    D --> E[答案生成]
    E --> F[事实性评测]
    F --> G{是否通过}
    G -->|是| H[返回答案]
    G -->|否| I[人工复核或重新生成]

这里,查询改写和权限检查都只依赖原始问题,因此可以并发;检索必须等待二者;答案生成必须等待检索;评测必须等待答案。

一个节点 vv 可执行的必要条件是:

upred(v),state(u)=Succeeded\forall u\in pred(v),\quad state(u)=Succeeded

其中 pred(v)pred(v)vv 的直接前驱集合。实际系统还要定义 SkippedFailedCancelled 等终态,因为“前驱结束”不等于“前驱成功”。

2. DAG 如何产生并发

假设任务耗时如下:

任务 耗时 依赖
查询改写 1 秒
权限检查 0.2 秒
检索 2 秒 查询改写、权限检查
生成 3 秒 检索
评测 1 秒 生成

串行执行需要:

1+0.2+2+3+1=7.2 秒1+0.2+2+3+1=7.2\text{ 秒}

但查询改写和权限检查可以并发,所需时间变为:

max(1,0.2)+2+3+1=7 秒\max(1,0.2)+2+3+1=7\text{ 秒}

并发只减少了独立分支的等待,不会消除依赖链。更一般地,DAG 的理论最短完成时间受最长路径限制:

Tcritical=maxpPaths(G)vpduration(v)T_{\text{critical}}=\max_{p\in Paths(G)}\sum_{v\in p} duration(v)

这称为关键路径长度。实际完成时间还会增加调度、队列、网络、限流和重试开销。

3. “看起来独立”不等于真正独立

两个任务可以并发,至少要同时满足:

  1. 输入互不依赖;
  2. 不读写相互冲突的可变状态;
  3. 资源限制允许并发;
  4. 失败时可以分别处理;
  5. 并发顺序不影响业务语义。

例如“生成答案”和“扣减账户余额”不能仅因为代码上没有参数依赖就并发。它们可能共享账户状态,且扣款具有外部副作用。

一个实用的依赖分类是:

  • 数据依赖:任务 B 需要任务 A 的输出;
  • 控制依赖:只有 A 满足条件,B 才能运行;
  • 资源依赖:同一 GPU、租户、文件或配额不能同时使用;
  • 一致性依赖:B 必须读取 A 已提交的版本;
  • 副作用依赖:B 必须在 A 的外部动作完成后执行。

只把显式函数参数当作依赖,会漏掉后四类依赖。

4. DAG 的边界:循环不是 DAG

LLM Agent 常见如下循环:

规划 -> 调用工具 -> 观察结果 -> 修正计划 -> 调用工具 -> ...

如果循环次数不预先展开,它不是 DAG,而是状态机或事件驱动过程。可以把最多三轮循环静态展开成 DAG,但这会产生固定上限,并不能表达任意终止条件。

循环应显式定义:

  • 当前轮次;
  • 最大轮数;
  • 每轮输入;
  • 终止条件;
  • 无进展检测;
  • 单轮和总预算;
  • 工具错误处理;
  • 人工接管条件。

例如,若每一轮都生成相同查询,系统不应无限重试。可以定义:

terminate=successroundRmaxcostBno_progressKterminate = success \lor round \geq R_{\max} \lor cost \geq B \lor no\_progress \geq K

这里 RmaxR_{\max} 是最大轮数,BB 是预算,KK 是连续无进展轮数。

三、事件:让工作流从“函数调用”变成“可恢复系统”

1. 命令、事件和状态不是同一个东西

三者应明确区分:

  • 命令(command):希望系统做什么,例如 RunRetrieval
  • 事件(event):已经发生了什么,例如 RetrievalSucceeded
  • 状态(state):系统当前认定的事实,例如任务状态为 Succeeded

命令可能丢失、重复或被拒绝;事件表示一个已经发生的事实,通常应追加写入;状态是根据事件或事务结果得到的当前视图。

一个事件至少应包含:

{
  "event_id": "evt-8f1",
  "event_type": "TaskSucceeded",
  "workflow_id": "wf-123",
  "task_id": "retrieve",
  "attempt": 2,
  "occurred_at": "2025-01-01T12:00:00Z",
  "payload": {
    "artifact_uri": "s3://bucket/runs/wf-123/retrieve.json",
    "input_digest": "sha256:...",
    "output_digest": "sha256:..."
  },
  "schema_version": 1
}

event_id 用于去重,attempt 用于区分执行尝试,input_digestoutput_digest 用于审计和复现。时间戳不能替代版本号,因为时钟可能漂移且事件可能乱序到达。

2. 事件驱动执行路径

典型流程如下:

sequenceDiagram
    participant C as Client
    participant O as Orchestrator
    participant DB as State Store
    participant Q as Queue
    participant W as Worker
    participant T as Tool

    C->>O: StartWorkflow
    O->>DB: 写入 WorkflowStarted
    O->>Q: 投递可执行 Task
    Q->>W: 领取 Task
    W->>T: 调用工具或模型
    T-->>W: 返回结果
    W->>DB: 事务写入 TaskSucceeded
    DB-->>O: 发布状态变更
    O->>DB: 计算后继节点
    O->>Q: 投递新的 Task

这里最重要的不是消息队列本身,而是“任务领取”和“状态提交”的边界。Worker 进程可能在工具调用成功后、写入 TaskSucceeded 前崩溃。恢复后,编排器只能看到任务没有成功,因而可能再次投递。

所以,事件驱动系统必须接受**至少一次执行(at-least-once)**的现实,并通过幂等设计处理重复。

3. 事件日志与当前状态

有两种常见方式:

方式 A:状态表为主,事件表用于审计

每次状态变更在同一事务中:

  1. 更新任务当前状态;
  2. 插入事件;
  3. 提交事务;
  4. 再异步发布消息。

优点是查询简单、恢复快。缺点是需要保证数据库提交和消息发布之间的一致性。

方式 B:事件溯源

只追加事件:

WorkflowStarted
TaskScheduled
TaskStarted
TaskSucceeded
TaskStarted
TaskFailed

当前状态通过重放事件得到。优点是审计和历史完整;缺点是事件模式演进、重放时间和查询复杂度更高。生产系统通常会保存快照,避免每次从第一条事件开始重放。

事件溯源不自动解决一致性。若事件本身不完整、顺序不明确或消费者处理非幂等,重放仍会产生错误。

4. 事件顺序与乱序

同一个任务的事件通常需要单调版本:

versionn+1=versionn+1version_{n+1}=version_n+1

消费者收到版本 5 时,如果本地只有版本 3,应暂缓应用,而不能假设版本 4 不重要。跨不同任务的事件则未必存在全局顺序,系统不应依赖消息到达顺序推导业务因果。

例如:

TaskA.Succeeded
TaskB.Succeeded

如果 A 和 B 无依赖,谁先到达都不影响结果;如果 B 的处理需要 A 的输出,则应由 DAG 依赖或显式版本条件保证,而不是依赖队列碰巧按顺序投递。

四、并发:不仅是 asyncio.gather

1. 并发、并行和异步的区别

  • 并发:多个任务在时间上交错推进;
  • 并行:多个任务同时占用不同 CPU、GPU 或执行单元;
  • 异步:调用方不必阻塞等待结果。

调用多个远程模型通常是 I/O 并发,不代表本机 CPU 并行。并发度过高会触发模型服务限流、GPU 显存不足、数据库连接耗尽或租户配额超限。

2. 资源约束下的可执行集合

设某一时刻可执行节点集合为 RR,但资源约束为:

vRcpuvCPUmax\sum_{v\in R} cpu_v \leq CPU_{\max}

vRgpuvGPUmax\sum_{v\in R} gpu_v \leq GPU_{\max}

vRcost_ratevBudgetRate\sum_{v\in R} cost\_rate_v \leq BudgetRate

因此,拓扑上可并发的任务不一定都应立即启动。调度器还需要考虑优先级、租户公平性、截止时间和成本预算。

3. 分支汇合必须定义缺失语义

设节点 merge 依赖三个分支:

A ─┐
B ─┼─> merge
C ─┘

如果 B 失败,merge 有至少四种不同语义:

  1. 全成功才运行:B 失败则整个流程失败;
  2. 允许部分成功:merge 接收 A、C 和 B 的错误;
  3. 失败分支降级:用缓存或默认值替代 B;
  4. 动态重算:重新选择一个替代 Agent。

必须把这种语义写入工作流定义,否则不同 Worker 可能对同一状态作出不同判断。

4. 竞态条件示例

两个 Agent 同时读取余额 100,各自扣除 80:

Agent A 读取 100
Agent B 读取 100
Agent A 写入 20
Agent B 写入 20

最终余额是 20,但正确结果应当拒绝至少一个请求。问题不是模型判断错误,而是“读取—计算—写入”不是原子操作。

可用以下方式修复:

  • 数据库事务和行锁;
  • 乐观锁:UPDATE ... WHERE version = old_version
  • 原子条件更新:UPDATE account SET balance = balance - 80 WHERE balance >= 80
  • 单租户串行队列;
  • 将扣款交给具备幂等和事务语义的专用服务。

让 LLM 输出“我已经扣款”不能产生事务性。模型输出只是数据,真正的状态变更必须由受控工具执行。

五、重试:处理暂时性故障,而不是重复犯错

1. 什么可以重试

适合重试的错误通常具有暂时性:

  • 网络连接失败;
  • 服务端 429 限流;
  • 5xx 短暂故障;
  • Worker 租约过期;
  • 依赖服务正在切换。

通常不应直接重试:

  • 参数校验失败;
  • 权限拒绝;
  • 资源不存在;
  • 明确的业务规则拒绝;
  • 模型输出违反不可重试的安全策略。

模型返回低质量答案也不自动意味着“重试同一个请求”。应先判断是否需要换 Prompt、增加检索上下文、切换模型或转人工。

2. 指数退避与抖动

常用的等待时间为:

dk=min(dmax,d02k)+Jd_k=\min(d_{\max},d_0\cdot 2^k)+J

其中:

  • kk 是已失败次数;
  • d0d_0 是初始等待;
  • dmaxd_{\max} 是上限;
  • JJ 是随机抖动。

抖动避免大量 Worker 同时失败后在同一时刻再次请求,形成惊群。

重试次数不能只由单任务决定。一个任务的第 5 次重试可能已经超出整个工作流的截止时间和成本预算。因此应同时设置:

  • 单任务最大尝试次数;
  • 工作流总截止时间;
  • 工作流总 token 和金额预算;
  • 同一错误的熔断阈值;
  • 同一输入的全局去重策略。

3. 重试与幂等的关系

如果一个操作是纯函数:

f(x)=yf(x)=y

重复执行通常没有副作用。但现实中的工具常常是:

发送邮件
创建订单
扣款
发布模型
写入外部 CRM

对这些操作重试可能重复产生副作用。解决方法是向下游传递幂等键:

Idempotency-Key = workflow_id + task_id + business_operation

下游服务需要保存:

key -> result

收到相同 key 时:

  • 若之前已成功,返回原结果;
  • 若之前正在处理,返回处理中或等待;
  • 若之前失败,要根据失败语义决定是否允许复用或重新执行。

仅仅在客户端生成幂等键还不够;下游必须真正识别并持久化它。

4. 超时不等于失败

客户端超时只表示“调用方没有在期限内得到响应”,不表示远程操作没有完成。

客户端发送创建订单
服务端已创建订单
网络响应丢失
客户端超时
客户端重试

如果没有幂等键,可能创建两个订单。更严谨的流程是:

  1. 超时后把本地任务标为 UnknownPendingConfirmation
  2. 用幂等键查询远程状态;
  3. 只有确认未创建时才重新发起;
  4. 无法确认时转人工或进入补偿流程。

许多系统错误地把所有超时都转为 Failed,这是重复副作用的常见根源。

六、一致性:模型上下文、工作流状态和外部系统必须分层处理

“一致性”不是一个单一属性,至少包括三层。

1. 编排状态一致性

编排器需要确保不会同时把同一任务错误地推进到两个互相冲突的状态。例如:

Pending -> Running -> Succeeded
Pending -> Cancelled

而以下转换通常不应直接发生:

Succeeded -> Running

除非明确表示重新执行,并增加新的执行版本。

可以为每个任务保存版本号:

UPDATE tasks
SET status = 'succeeded',
    version = version + 1,
    output_uri = :output_uri
WHERE workflow_id = :workflow_id
  AND task_id = :task_id
  AND status = 'running'
  AND version = :expected_version;

若影响行数为 0,说明状态已经被其他执行者修改,当前 Worker 不应继续发布成功事件。

2. 数据一致性

模型、检索索引、特征库、评测集和答案生成可能看到不同版本的数据。

例如:

文档库版本 v10
向量索引版本 v9
答案生成读取索引 v9
评测读取文档库 v10

此时评测结果不能直接与生成结果比较,因为输入语料已经不一致。生产运行应记录:

  • 数据集版本;
  • 文档快照或内容摘要;
  • Embedding 模型版本;
  • 索引构建版本;
  • 生成模型版本;
  • Prompt 模板版本;
  • 工具和代码版本。

可将运行输入表示为:

InputFingerprint=Hash(dataset_version,index_version,model_version,prompt_version,parameters)InputFingerprint = Hash(dataset\_version, index\_version, model\_version, prompt\_version, parameters)

相同指纹表示“有机会复现”,但不保证 LLM 输出逐 token 相同,因为服务端采样、系统提示和底层实现仍可能变化。

3. 外部副作用一致性

工作流数据库和外部服务通常不是同一个事务,因此不能假设跨服务原子提交。

例如:

数据库写入“邮件已发送”
但邮件服务调用失败

或:

邮件已经发送
数据库提交“已发送”之前 Worker 崩溃

常见解决方案是 Outbox 模式

  1. 在同一数据库事务中写入业务状态和待发送 Outbox 记录;
  2. 独立发布器读取 Outbox;
  3. 发布成功后标记 Outbox;
  4. 下游使用幂等键处理重复消息。

Outbox 仍不是“恰好一次执行”。它通常提供:

  • 数据库内状态和待发布消息的一致提交;
  • 消息至少一次投递;
  • 通过下游幂等实现业务上的有效一次。

“恰好一次”往往只能在有限边界内成立,例如单个数据库事务内的状态更新。跨网络、跨数据库和外部副作用的全局 exactly-once 通常不可直接假设。

七、多 Agent 编排:协作关系必须显式化

1. 常见拓扑

多 Agent 系统至少有几种结构:

顺序委派

协调 Agent -> 检索 Agent -> 分析 Agent -> 写作 Agent

适合依赖强、结果需要逐步加工的流程。

并行专家

             -> 安全 Agent -
协调 Agent  -> 事实 Agent  -> 汇总 Agent
             -> 风险 Agent -

适合多个独立视角,但汇总 Agent 必须处理冲突和证据来源。

层级管理

总控 Agent
  ├── 数据任务 Agent
  ├── 模型任务 Agent
  └── 发布任务 Agent

总控 Agent 分配目标,子 Agent 在边界内执行。层级越深,状态传播、权限和调试越复杂。

竞争与裁决

多个 Agent 独立生成结果,再由裁决器选择。它可能提高鲁棒性,但会增加 token、延迟和一致性成本,且“多数投票”不等价于事实正确。

2. Agent 之间传什么

不要默认把完整对话历史复制给所有 Agent。应区分:

  • 共享工作流状态:任务 ID、数据版本、审批状态;
  • 任务输入:当前 Agent 必需的字段;
  • 证据引用:文档 URI、检索结果 ID、哈希;
  • 模型上下文:只传当前推理所需的内容;
  • 工具结果:原始响应、规范化响应和错误分类。

大上下文复制会增加 token 成本,也可能造成越权信息泄露。更好的方式是传递受权限控制的 artifact 引用,并在读取时重新校验访问权。

3. Handoff 与共享状态的边界

Agent 转交任务时,至少要定义:

  • 转交目标;
  • 转交原因;
  • 结构化输入;
  • 已完成动作;
  • 未完成动作;
  • 证据和引用;
  • 允许的工具;
  • 截止时间和预算;
  • 终止或退回条件。

如果只传一句自然语言“请继续处理”,下一个 Agent 无法可靠区分已完成和未完成的动作,容易重复调用工具。

一个安全的转交对象可以是:

{
  "case_id": "case-42",
  "objective": "判断退款是否符合政策",
  "completed": [
    {"action": "load_order", "artifact": "orders/42.json"}
  ],
  "pending": ["check_policy", "request_approval_if_needed"],
  "evidence": [
    {"uri": "orders/42.json", "sha256": "..."}
  ],
  "constraints": {
    "max_tool_calls": 4,
    "requires_human_approval_for": ["refund_over_1000"]
  }
}

4. 终止和人工确认

Agent 循环必须有机器可验证的终止条件,不能只依赖模型说“任务完成”。例如:

  • 必需字段全部存在;
  • 工具调用次数达到上限;
  • 输出通过结构化 Schema 校验;
  • 评测分数达到阈值;
  • 风险等级要求人工批准;
  • 到达截止时间或预算上限。

人工确认是状态转换,不是 Prompt 中的一句“请确认”:

AwaitingApproval -> Approved -> Publish
AwaitingApproval -> Rejected -> Cancelled
AwaitingApproval -> Expired -> Escalated

暂停期间应持久化完整上下文、审批人、审批时间、审批依据和上下文版本。恢复时重新检查权限和数据版本,不能盲目使用几小时前的授权结果。

八、MCP 在编排中的位置

Model Context Protocol(MCP)定义了模型应用与外部能力之间的协议边界。其核心角色通常包括:

  • Host:承载 AI 应用或 Agent;
  • Client:Host 内与某个 MCP Server 建立协议连接的组件;
  • Server:提供工具、资源或提示模板的一方。

MCP 使用基于 JSON-RPC 的消息交互,并通过初始化和能力协商了解双方支持的能力。工具通常具有名称、描述和输入 Schema;资源用于暴露可读取的数据;提示模板用于提供可复用的交互模板。具体字段、生命周期和能力范围应以当前 MCP Specification 为准,因为协议会演进。

MCP 解决的是“如何以标准协议发现和调用外部上下文与工具”,它本身不规定:

  • DAG 如何调度;
  • 任务是否至少一次执行;
  • 重试次数和退避策略;
  • 工作流状态如何持久化;
  • 多 Agent 如何仲裁;
  • 外部副作用如何实现幂等;
  • 业务审批如何完成。

因此,正确的分层是:

工作流编排器
  ├── DAG / 状态机 / 事件 / 重试 / 预算
  └── Agent Runtime
        └── MCP Client
              └── MCP Server
                    └── 数据库、API、文件或业务系统

调用 MCP 工具时,编排器仍应记录:

  • 工具名称和版本;
  • 输入 Schema 验证结果;
  • 工作流和任务幂等键;
  • 调用开始、结束和错误事件;
  • 资源访问权限;
  • 返回内容的摘要或 artifact 地址;
  • 是否发生了外部副作用。

MCP Server 不应仅因为客户端声称“这是内部 Agent”就信任请求。工具服务必须独立执行认证、授权、参数校验、租户隔离和审计。工具描述也不是安全策略;模型可能误解描述,恶意内容也可能诱导 Agent 调用高风险工具。

OpenAI 的 Agents 相关指南所覆盖的 Agent、工具、交接、护栏和追踪等概念,可以作为 Agent Runtime 的设计参考。但具体 SDK 方法、参数和可用能力会随版本变化,生产代码应以当前 SDK 文档和类型定义为准,不应把某个示例接口当作跨版本稳定的编排协议。无论使用哪种 SDK,DAG 状态、幂等键、事件持久化和人工审批都仍应由系统层明确设计。

九、一个可运行的最小 DAG 编排器

下面的 Python 示例只使用标准库,演示:

  • DAG 依赖;
  • 两个独立任务并发;
  • 任务级重试;
  • 失败后不执行后继节点;
  • 结构化事件日志。

它不是生产级持久化编排器,但可以直接运行以观察核心机制。

import asyncio
import random
from collections import defaultdict

GRAPH = {
    "rewrite": [],
    "auth": [],
    "retrieve": ["rewrite", "auth"],
    "generate": ["retrieve"],
    "evaluate": ["generate"],
}

async def run_task(name, inputs, attempt):
    await asyncio.sleep({
        "rewrite": 0.3,
        "auth": 0.1,
        "retrieve": 0.4,
        "generate": 0.5,
        "evaluate": 0.2,
    }[name])

    # 模拟一次暂时性故障:retrieve 的第一次尝试失败
    if name == "retrieve" and attempt == 1:
        raise ConnectionError("temporary retrieval failure")

    if name == "auth":
        return {"authorized": True}
    if name == "rewrite":
        return {"query": inputs["question"].lower()}
    if name == "retrieve":
        return {"documents": ["doc-1", "doc-2"]}
    if name == "generate":
        return {"answer": "基于 doc-1 和 doc-2 生成的答案"}
    if name == "evaluate":
        return {"passed": True}

    raise ValueError(f"unknown task: {name}")

async def execute(name, results, max_attempts=3):
    inputs = {"question": "How does retry work?"}
    for dependency in GRAPH[name]:
        inputs[dependency] = results[dependency]

    for attempt in range(1, max_attempts + 1):
        print({
            "event": "TaskStarted",
            "task": name,
            "attempt": attempt,
        })
        try:
            output = await run_task(name, inputs, attempt)
            print({
                "event": "TaskSucceeded",
                "task": name,
                "attempt": attempt,
                "output": output,
            })
            return output
        except ConnectionError as exc:
            print({
                "event": "TaskRetryableFailure",
                "task": name,
                "attempt": attempt,
                "error": str(exc),
            })
            if attempt == max_attempts:
                raise
            await asyncio.sleep(0.1 * 2 ** (attempt - 1)
                               + random.uniform(0, 0.05))

async def main():
    results = {}
    pending = set(GRAPH)
    completed = set()

    while pending:
        ready = [
            name for name in pending
            if set(GRAPH[name]).issubset(completed)
        ]
        if not ready:
            raise RuntimeError("DAG has a cycle or an unreachable task")

        # 当前 ready 集合中的任务并发执行
        batch = await asyncio.gather(
            *(execute(name, results) for name in ready),
            return_exceptions=True,
        )

        for name, result in zip(ready, batch):
            if isinstance(result, Exception):
                print({
                    "event": "WorkflowFailed",
                    "task": name,
                    "error": repr(result),
                })
                return
            results[name] = result
            completed.add(name)
            pending.remove(name)

    print({"event": "WorkflowSucceeded", "results": results})

if __name__ == "__main__":
    asyncio.run(main())

运行时,rewriteauth 会同时开始;retrieve 等待二者成功;它第一次调用抛出可重试的 ConnectionError,经过退避后第二次成功;最后才会运行 generateevaluate

这个示例有几个故意保留的限制:

  1. results 在内存中,进程崩溃后全部丢失;
  2. 没有 Worker 租约,两个进程可能同时执行同一任务;
  3. 任务结果写入和事件发布不是事务性的;
  4. 没有幂等键,不能安全处理外部副作用;
  5. auth 的结果只是内存值,不能替代真实权限系统;
  6. evaluate 的通过条件固定,不能表示复杂的人工确认。

要进入生产环境,至少需要把 results、任务状态、尝试次数和事件日志持久化,并使用条件更新防止重复领取。

十、生产状态机应比 success/fail 更细

一个任务的状态可以设计为:

Pending
  -> Scheduled
  -> Running
  -> Succeeded
  -> Failed
  -> Retrying
  -> TimedOut
  -> Cancelled
  -> AwaitingApproval
  -> Unknown

关键状态的语义如下:

  • Pending:依赖尚未满足;
  • Scheduled:已进入队列,但还没有 Worker 执行;
  • Running:持有执行租约;
  • Retrying:本次尝试失败,等待下一次;
  • TimedOut:超过执行期限,结果可能未知;
  • Unknown:无法判断外部副作用是否已发生;
  • AwaitingApproval:工作流主动暂停;
  • Succeeded:输出已持久化并满足完成条件;
  • Failed:已确定不能通过当前策略完成;
  • Cancelled:被用户、策略或上游失败取消。

Unknown 很重要。它表达的是认知上的不确定,而不是业务失败。对于发送邮件、创建支付、发布模型等操作,超时后直接进入 Failed 会诱发重复执行。

Worker 领取任务时可以采用租约:

UPDATE tasks
SET status = 'running',
    lease_owner = :worker_id,
    lease_expire_at = now() + interval '60 seconds',
    attempt = attempt + 1
WHERE task_id = :task_id
  AND status IN ('scheduled', 'retrying')

租约过期后,调度器可以重新投递任务。但旧 Worker 可能在网络分区后恢复,因此提交结果时仍必须带上执行版本或租约令牌,不能只检查 task_id

十一、失败传播和补偿

DAG 中的失败传播不能只写成“某节点失败,整个流程失败”。应为每条边定义策略:

  • fail_fast:前驱失败立即取消后继;
  • continue_on_error:后继接收错误对象;
  • fallback:使用备用节点或缓存;
  • manual_review:进入人工审批;
  • compensate:执行反向业务动作。

补偿不是回滚。数据库事务可以回滚未提交的更新,但外部邮件、已发出的 HTTP 请求和已经训练完成的 GPU 作业不能真正回滚。补偿只能执行另一个动作,例如:

创建资源成功 -> 删除资源
扣款成功 -> 发起退款
发布成功 -> 切换回旧版本

补偿动作也可能失败,因此它需要自己的状态、重试、幂等键和人工接管路径。

一个典型的发布流程可以是:

训练成功
  -> 离线评测通过
  -> 安全扫描通过
  -> 人工批准
  -> 灰度发布
  -> 在线指标观察
  -> 全量发布

若全量发布后指标异常,回滚到旧模型通常是切换流量指针,而不是删除新模型。模型文件、数据快照、评测报告和发布记录应保持可追溯。

十二、模型、数据、评测、权限和成本要进入同一状态图

AI 生产系统的一个常见错误,是只持久化“模型输出”,不持久化产生输出的条件。

模型版本

至少记录:

  • 模型标识和版本;
  • 推理参数;
  • Prompt 或系统指令版本;
  • 工具描述版本;
  • 输出 Schema;
  • 安全策略版本。

数据版本

至少记录:

  • 训练或评测数据集版本;
  • RAG 文档快照;
  • 索引版本;
  • Embedding 模型版本;
  • 数据权限过滤条件。

评测状态

评测不应只是一个浮点数。应同时记录:

{
  "metric": {
    "factuality": 0.92,
    "format_valid": 1.0
  },
  "dataset_version": "eval-v17",
  "threshold_version": "policy-v4",
  "passed": true
}

评测通过表示“在指定数据集、指标定义和阈值下通过”,不是绝对正确。

权限

权限检查应尽量靠近工具执行点。工作流开始时的授权不能自动覆盖整个长流程,因为:

  • 用户权限可能被撤销;
  • Agent 可能转交给另一个身份;
  • 工具参数可能扩大访问范围;
  • 人工批准可能只允许某个版本或某个租户。

权限决策应绑定主体、资源、动作、租户和上下文版本,并写入审计日志。

成本

一次 Agent 运行的成本可以粗略写成:

C=i(input_tokensipiin+output_tokensipiout)+jtool_costj+kinfra_costkC = \sum_i (input\_tokens_i \cdot p^{in}_i + output\_tokens_i \cdot p^{out}_i) +\sum_j tool\_cost_j +\sum_k infra\_cost_k

其中还应考虑重试和并发分支。若一个节点平均执行 aa 次,预计成本并不是单次成本,而是:

E[Ctask]=k=1aP(达到第 k 次)CkE[C_{task}]=\sum_{k=1}^{a}P(\text{达到第 }k\text{ 次})\cdot C_k

因此,成本预算应是可执行约束,而不是运行结束后的报表。达到预算后可以停止低优先级分支、切换小模型、缩短上下文或请求人工处理。

十三、诊断:从事件链而不是最后一条错误开始

一次失败运行至少要能按以下维度关联:

trace_id
workflow_id
task_id
attempt
worker_id
model_call_id
tool_call_id
artifact_id
tenant_id

诊断时按这条路径检查:

  1. 工作流是否创建成功;
  2. DAG 是否正确计算出可执行节点;
  3. 任务是否进入队列;
  4. Worker 是否成功领取租约;
  5. 模型或工具请求是否发出;
  6. 超时发生在客户端、网关还是下游;
  7. 外部副作用是否可能已经发生;
  8. 成功结果是否持久化;
  9. 后继节点为什么没有被调度;
  10. 是否因权限、预算或人工审批暂停。

常见症状和原因并不一一对应:

症状 可能原因
任务反复执行 租约过期、提交版本过旧、缺少幂等键
后继节点永远不运行 前驱状态未成功、事件乱序、依赖定义错误
工作流显示失败但外部订单存在 超时后错误分类为失败
评测结果不可复现 数据、索引、Prompt 或模型版本未固定
并发一提高就大量 429 未按租户或模型服务设置限流
Agent 无限循环 没有轮数、预算、无进展和终止条件
两个 Agent 覆盖彼此结果 共享状态没有版本检查或冲突策略

日志应记录输入和输出的摘要、Schema 校验结果、token 用量和延迟,但不应无条件记录敏感 Prompt、个人数据或完整工具响应。调试需要可见性,合规要求最小化暴露,二者应通过脱敏、分级访问和短期调试采样平衡。

十四、常见误解与边界

误解一:DAG 能描述所有 Agent 流程

不能。固定依赖适合 DAG;动态循环、条件转交和不确定终止更适合状态机或事件驱动模型。二者可以组合:外层用 DAG 管理阶段,阶段内部由 Agent 状态机运行。

误解二:消息队列保证不重复

通常不保证。即使队列提供某种消息去重,也无法覆盖 Worker 在外部副作用完成后崩溃的情况。业务幂等必须在任务和下游服务两侧设计。

误解三:把温度设为零就能完全复现

不能。模型服务版本、系统提示、检索顺序、工具返回、并发时序和底层实现都可能变化。低随机性有助于稳定,但复现仍需要锁定输入和依赖版本。

误解四:多个 Agent 互相审查就一定更可靠

不一定。若所有 Agent 使用相同错误资料、相同错误假设或相同提示注入,投票只会放大共同错误。可靠性提升来自独立证据、结构化校验、规则约束和可验证工具,而不是 Agent 数量本身。

误解五:工具 Schema 就是安全边界

不是。Schema 主要约束数据形状;认证、授权、资源范围、速率限制和副作用确认仍必须由工具服务执行。

误解六:重试次数越多越可靠

重试只能提高暂时性故障下的成功概率,同时也增加延迟、成本和重复副作用风险。对确定性错误重试没有收益,对未知状态操作重试可能有害。

十五、设计检查的核心问题

设计一个 AI 工作流或多 Agent 系统时,必须能回答以下问题:

  1. 每个节点的输入、输出和副作用是什么?
  2. 哪些边是数据依赖,哪些是资源或一致性依赖?
  3. 哪些节点可并发,最大并发度由什么资源约束?
  4. 循环的最大轮数、预算和无进展条件是什么?
  5. 每个错误是永久失败、暂时失败还是未知状态?
  6. 重试时使用什么幂等键?
  7. Worker 崩溃在每个时间窗口内会发生什么?
  8. 状态提交和事件发布如何保持一致?
  9. 外部副作用发生但本地没有记录时如何确认或补偿?
  10. 模型、Prompt、数据、索引、评测和权限版本如何关联?
  11. 多 Agent 转交时传递哪些结构化状态?
  12. 人工审批暂停和恢复时如何重新验证授权?
  13. 成本、延迟和 token 预算在哪些状态转换处强制执行?
  14. 能否从一个 workflow_id 重建完整执行路径?

当这些问题都有明确答案时,DAG、事件、并发、重试和一致性就不再是彼此孤立的术语,而会共同形成一套可恢复的执行语义:

可执行性=依赖满足权限有效资源可用预算未超\text{可执行性} = \text{依赖满足} \land \text{权限有效} \land \text{资源可用} \land \text{预算未超}

可靠完成=状态持久化副作用可确认重复执行可接受终止条件可验证\text{可靠完成} = \text{状态持久化} \land \text{副作用可确认} \land \text{重复执行可接受} \land \text{终止条件可验证}

这也是 AI Agent 从演示代码进入生产系统时必须跨越的边界:模型负责生成候选决策,编排器负责控制状态和故障,工具服务负责执行并保护副作用,事件和版本负责说明系统究竟发生了什么。


系列导航与关联阅读

官方资料

本文依据研究论文、标准组织与主流框架官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。