跳转至

深入专题:工作流、恢复与副作用

这篇解决一个生产问题:

一个 Agent 任务运行了 20 分钟,执行到一半服务重启、模型超时或等待人工审批,系统怎样知道从哪里继续,并且不重复产生副作用?

先区分 Workflow 和 Agent

Workflow 的下一步由程序规则决定:

解析 → 校验 → 查询 → 生成 → 审批 → 发布

Agent 的下一步可能由模型根据上下文决定:

模型观察当前状态
→ 选择工具或结束

二者不是竞争关系。生产系统常见结构是:

确定性 Workflow
├── 输入校验
├── 权限
├── 数据准备
├── Agent 决策节点
├── 人工审批
├── 副作用执行
└── 审计和终态

确定的事情交给程序,不确定的局部语义判断交给模型。

为什么长任务不能只存在进程内存

假设 Python 变量保存:

state = {
    "task_id": "T-100",
    "current_step": "analyze",
    "results": [...]
}

进程崩溃、容器迁移或发布重启后,这些状态就丢失。用户只看到任务永远处于“处理中”。

生产任务至少需要持久化:

{
  "task_id": "T-100",
  "workflow_version": "diagnosis-v3",
  "status": "running",
  "current_node": "query_alarm",
  "input": {"device_id": "D-1"},
  "state": {"completed_queries": ["device_info"]},
  "attempt": 2,
  "owner": "worker-7",
  "lease_until": "...",
  "created_at": "...",
  "updated_at": "..."
}

状态里存原始事实,不存拼好的 Prompt

不推荐:

{
  "context": "设备 D-1 是……以下是一大段已经格式化的提示词"
}

推荐:

{
  "device": {"id": "D-1", "model": "M-2"},
  "online_events_ref": "object://task/T-100/online.json",
  "alarm_summary": {"count": 3},
  "analysis": null
}

原因:

  • 原始事实可以被不同节点复用;
  • Prompt 模板升级时可重新渲染;
  • 大结果可以存引用,避免 Checkpoint 膨胀;
  • 调试时能区分数据错误和提示词错误;
  • 敏感字段更容易分级处理。

节点边界就是恢复边界

一个节点做得太大:

查询 4 个系统
+ 调模型
+ 创建工单
+ 发通知

任何一步失败都很难判断哪些已完成。

更合理:

load_device
query_online_history
query_alarms
analyze
request_approval
create_ticket
notify

每个节点:

  • 输入和输出明确;
  • 可以单独记录耗时和错误;
  • 有清晰的重试策略;
  • 副作用能分配独立幂等键;
  • 可以从最近成功边界恢复。

节点也不能无限细,否则 Checkpoint 数量、调度开销和理解成本上升。边界通常放在:

  • 外部调用前后;
  • 重要业务状态变化处;
  • 人工等待处;
  • 不同重试策略之间;
  • 不同权限或副作用等级之间。

Checkpoint 保存的是什么

Checkpoint 不是简单的“当前步骤字符串”。它通常应包含:

  • 当前工作流状态快照;
  • 已完成节点;
  • 下一批可运行节点;
  • 节点写入结果;
  • 工作流版本;
  • 与本次运行关联的稳定 ID。

恢复时不是重新问模型“你刚才做到哪了”,而是运行时读取持久状态,按确定规则恢复。

至少一次执行意味着节点可能重跑

大多数可靠任务系统更容易提供:

某个节点至少会执行一次,但在故障窗口中可能执行多次。

例子:

Worker 执行 create_ticket
→ 工单创建成功
→ Worker 在写 Checkpoint 前崩溃

恢复后,系统只看见 create_ticket 尚未完成,于是再次执行。

所以:

可恢复执行 ≠ 自动获得恰好一次副作用

副作用节点必须自己幂等。

副作用应使用稳定操作 ID

operation_id =
workflow_run_id + node_name + business_target + version

例如:

T-100:create-ticket:incident-9:v1

服务端使用数据库唯一约束保存:

operation_id
request_hash
status
result

同一节点重放时返回已有工单结果。

人工审批为什么需要真正暂停

错误做法:

循环查询数据库:
每 5 秒检查一次是否批准

问题:

  • 占用 Worker;
  • 大量无意义查询;
  • 服务重启后状态难管理;
  • 几天没人审批就持续浪费资源。

正确模式:

生成审批请求
→ 保存 Checkpoint
→ 任务进入 waiting_for_approval
→ 释放 Worker

用户批准
→ 发送恢复事件
→ 加载同一运行状态
→ 从审批后继续

等待一天不应占用一天线程。

恢复时节点可能从开头重放

以 LangGraph 的动态中断为例,恢复后包含中断的节点会从节点开头重新执行。因此:

def approval_node(state):
    create_audit_record()   # 恢复时可能再次执行
    approved = interrupt(...)

create_audit_record() 必须幂等,或者拆成独立节点:

create_approval_request
→ wait_for_approval
→ execute_action

框架的持久化和中断提供恢复能力,但不会替你自动设计业务副作用。

任务领取需要 Lease,不只是状态字段

多个 Worker 可能同时看见一个待执行任务:

Worker A:读取 pending
Worker B:读取 pending

如果只是分别更新为 running,两者可能都执行。

常见领取方式:

UPDATE tasks
SET owner = :worker,
    status = 'running',
    lease_until = :future,
    version = version + 1
WHERE id = :id
  AND status IN ('pending', 'retrying')
  AND version = :expected_version;

只有受影响行数为 1 的 Worker 获得执行权。

Lease 到期允许其他 Worker 接管,以处理原 Worker 崩溃。但接管之后仍要依赖幂等来防止前一个 Worker 实际仍在运行。

更严格的系统会使用 fencing token:

每次接管 version 增加
下游只接受最新 version 的写入

心跳解决什么,不能解决什么

Worker 执行长节点时定期刷新:

lease_until = now + 30s

它能说明 Worker 最近仍活跃,但不能证明:

  • 外部工具一定没卡住;
  • 业务结果一定正确;
  • 网络分区另一边的旧 Worker 已停止;
  • 同一个副作用没有发生两次。

心跳是调度活性机制,不是业务正确性机制。

超时要分层

一个 15 分钟的任务不能只有一个总超时。

HTTP 请求:30 秒
模型调用:60 秒
只读工具:10 秒
人工审批:7 天
整个任务:30 分钟执行时间,不含人工等待

否则可能出现:

  • 工具早已卡死,但任务总超时很长;
  • 人工审批期间被误判为任务超时;
  • HTTP 连接一直占用,用户无法关闭页面。

长任务通常采用:

POST /tasks → 立即返回 task_id
GET /tasks/{id} → 查询进度
WebSocket/SSE → 推送步骤事件

取消任务不是删除任务

用户点击取消时:

status = cancel_requested

运行节点在安全点检查:

  • 还没产生副作用:停止;
  • 正在执行不可中断操作:等待完成再停止;
  • 已产生可补偿副作用:进入补偿;
  • 已产生不可逆副作用:记录“部分完成”,人工处理。

不能简单删除数据库记录,因为:

  • Worker 可能仍在执行;
  • 审计链丢失;
  • 已产生的外部影响不会随记录删除而消失。

错误需要分类

类型 示例 流程处理
瞬时基础设施错误 网络超时、临时 503 有界重试
模型可修复错误 输出结构不合法 带错误反馈重试一次
用户可修复错误 缺设备、缺时间 暂停并询问
业务终态错误 无权限、工单已关闭 明确失败,不重试
未知程序错误 空指针、反序列化异常 失败并告警,保留现场

错误是工作流的一部分,而不只是日志里的一行堆栈。

工作流版本升级是隐藏难题

运行中的任务可能在 v1 节点等待审批,此时部署了 v2:

  • 节点改名;
  • State 字段变化;
  • 路由规则变化;
  • Prompt 或工具参数升级。

恢复时用新代码解释旧 Checkpoint,可能失败。

可选策略:

  1. 为工作流记录版本,让旧任务仍由旧版本完成;
  2. 编写显式状态迁移;
  3. 对无法迁移的旧任务人工终止并重新发起;
  4. 发布前统计所有暂停和运行中的旧版本任务。

不要假设“代码向后兼容”会自然发生。

完整推演:可恢复设备诊断与建单

状态

{
  "run_id": "RUN-8",
  "workflow_version": "diagnosis-v3",
  "device_id": "D-1001",
  "facts": {},
  "analysis": null,
  "proposed_action": null,
  "approval": null,
  "ticket": null,
  "status": "running"
}

时间线

09:00  创建 RUN-8
09:01  load_device 成功,Checkpoint C1
09:02  query_history 成功,Checkpoint C2
09:03  analyze 模型超时
09:03  退避后第二次 analyze 成功,Checkpoint C3
09:04  创建审批请求,Checkpoint C4
09:04  状态 waiting_for_approval,释放 Worker
11:20  用户批准
11:20  恢复 RUN-8
11:21  create_ticket 成功,但 Worker 在保存 C5 前崩溃
11:22  新 Worker 接管,再次调用相同 operation_id
11:22  工单服务返回第一次创建的 TICKET-77
11:22  保存 C5
11:23  notify 成功,任务 completed

这个例子证明了什么

  • Checkpoint 解决从哪里恢复;
  • Retry 解决瞬时失败;
  • Interrupt 解决长时间等待;
  • Lease 解决 Worker 接管;
  • 幂等解决副作用重放;
  • 状态机解决合法流转;
  • Trace 解决事后解释。

任何单独一个机制都不等于“可靠工作流”。

何时用 LangGraph,何时看 Temporal

场景 更关注的能力
模型节点、条件路由、消息状态、人工中断 LangGraph 的图状态和 Checkpoint
跨天业务流程、活动重试、定时器、长期可靠调度 Temporal 一类持久工作流引擎
简单固定流程、规模小 普通任务队列加数据库状态机也可能足够

不要因为岗位写了框架名称就全部堆进项目。先根据故障恢复、等待时长、团队能力和运维成本选择。

权威延伸资料

自测

  1. 为什么把 current_step 存进数据库仍不算完整 Checkpoint?
  2. 为什么节点越大,恢复越困难?节点是否越小越好?
  3. “框架支持持久化”为什么不能保证创建工单只发生一次?
  4. 人工审批等待时,为什么不应占用 Worker?
  5. Lease、心跳、幂等各自解决什么不同问题?
  6. 用户取消任务时,为什么不能直接删记录?
  7. 工作流代码升级后,暂停中的旧任务有什么风险?
参考答案要点
  1. 还需保存节点输入输出、已完成写入、下一步、版本和运行标识,否则无法确定性恢复。
  2. 大节点含多个故障窗口,无法判断已完成部分;过小节点又增加调度和状态成本,应以外部调用、副作用和策略边界划分。
  3. 写外部系统成功、写 Checkpoint 前崩溃时节点会重放,必须由副作用接口使用稳定幂等键。
  4. 审批可能持续数小时或数天,应持久化状态并释放计算资源,由审批事件唤醒。
  5. Lease 决定谁暂时拥有任务;心跳延长所有权并反映活性;幂等保证重复执行不重复产生业务效果。
  6. 运行中的 Worker 和外部副作用不会随记录删除而消失,且删除破坏审计。应请求取消并在安全点收敛。
  7. 新代码可能无法识别旧节点、State 或路由,需要版本固定、状态迁移或显式处理旧任务。

完成标准

  • 能为 6 个以上节点的任务设计持久化 State
  • 能指出每个节点的重试、暂停和终态策略
  • 能推演“副作用成功、Checkpoint 失败”的恢复过程
  • 能解释 Lease、Checkpoint、幂等和状态机的边界
  • 能设计人工审批和取消流程
  • 完成 可恢复 Agent Workflow 项目

继续深入:RAG、Eval 与线上治理