第 8 章 · 工作流引擎与 Human-in-the-Loop
Agent 的自主性是资产也是负债。工作流(DAG 引擎)是把确定性控制流从模型手里拿回来的基础设施;Human-in-the-Loop(HITL)是在关键决策点把控制权交还给人。两者合起来回答一个问题:自主性放到什么程度,在哪里设闸门。
8.1 为什么需要工作流引擎
考虑一个生产任务:"分析这份财报 → 生成摘要 → 风控审查 → 发邮件"。三种实现:
- 全自主 Agent:让模型自己决定调工具的顺序。风险:顺序不可控(没审查就发了邮件)、中途失败无法恢复、成本不可预测;
- 硬编码脚本:代码里写死四步。可靠,但每一步内部的 LLM 调用细节、重试、并行都要手写,改一个步骤 = 发一个版本;
- 工作流引擎:步骤与依赖声明式定义(图),执行、重试、状态、恢复、事件全部由引擎统一提供。LLM 调用只是图上的一种节点类型。
工作流不是 Agent 的对立面——生产系统的正确形态常常是:工作流定义骨架,Agent 填充其中自主性最强的一两个节点。骨架保证流程正确与可恢复,Agent 节点处理无法预先确定的部分。
8.2 最小 DAG 执行器
用 ~80 行实现一个真实可用的 DAG 执行器(理解引擎的最佳方式):
"""minimal_dag.py — 拓扑排序 + 状态机 + 重试的最小工作流引擎"""
import time, json
from collections import defaultdict
class Graph:
def __init__(self):
self.nodes, self.edges = {}, defaultdict(list) # edges[a] = [下游...]
def node(self, node_id, fn, retries=2):
self.nodes[node_id] = {"fn": fn, "retries": retries, "deps": set()}
return node_id
def edge(self, upstream, downstream):
self.edges[upstream].append(downstream)
self.nodes[downstream]["deps"].add(upstream) # 下游记录依赖
def validate(self):
"""环检测:Kahn 算法,有环直接拒绝执行"""
indeg = {n: len(d["deps"]) for n, d in self.nodes.items()}
queue = [n for n, d in indeg.items() if d == 0]
seen = 0
while queue:
n = queue.pop()
seen += 1
for m in self.edges[n]:
indeg[m] -= 1
if indeg[m] == 0:
queue.append(m)
if seen != len(self.nodes):
raise ValueError("图中存在环,拒绝执行")
class Executor:
def __init__(self, graph):
self.g = graph
self.state = {n: "pending" for n in graph.nodes} # pending/running/done/failed/skipped
self.outputs = {}
def ready_nodes(self):
"""依赖全部完成的 pending 节点 → 本轮可并发执行"""
return [n for n, s in self.state.items() if s == "pending"
and all(self.state[d] == "done" for d in self.g.nodes[n]["deps"])]
def run_node(self, node_id):
spec = self.g.nodes[node_id]
self.state[node_id] = "running"
for attempt in range(spec["retries"] + 1):
try:
self.outputs[node_id] = spec["fn"](
{d: self.outputs[d] for d in spec["deps"]}) # 注入上游输出
self.state[node_id] = "done"
return
except Exception as e:
if attempt == spec["retries"]:
self.state[node_id] = "failed"
for m in self.g.edges[node_id]: # 级联跳过下游
self.state[m] = "skipped"
raise
time.sleep(2 ** attempt) # 指数退避
def run_until_done(graph, executor, poll=0.05):
"""调度循环:生产实现改为事件驱动 + 并发池"""
while any(s == "pending" for s in executor.state.values()):
for n in executor.ready_nodes():
executor.run_node(n)
return executor.state, executor.outputs
这个骨架已经包含引擎的全部核心概念:图校验先于执行(环检测)、就绪集合(依赖满足才可运行)、节点级重试、失败级联语义(下游 skip 而非 fail,终态可区分"没跑"和"跑了失败")。生产引擎在此之上加:并发调度池、状态与输出持久化(每节点先落库再执行)、暂停/恢复、事件总线。
8.3 状态持久化与恢复
工作流的可恢复性来自一条纪律:状态变更先落库,执行才发生(write-ahead)。
节点状态的转移表:
pending --[依赖就绪]--> running --[成功]--> done
└--[重试耗尽]--> failed(下游 → skipped)
pending/running --[用户取消]--> cancelled
崩溃恢复的语义:进程重启后扫描所有 running 节点——无法确认是否已执行(可能执行到一半崩了)。两种处理:
- 节点声明幂等(同样输入重复执行结果一致)→ 直接重跑;
- 不幂等(发邮件、付款)→ 执行器需要"执行前意向日志"(intent log):先记"准备执行",再执行,再记"已执行"。恢复时对有 intent 但无结果的任务询问或走补偿。
这也是为什么生产工作流把"发送类"节点设计成两段式:先落"待发送"(幂等可查重),由独立的投递器去重后真正发送。
8.4 Human-in-the-Loop:审批是一等公民
HITL 不是 UI 上的确认弹窗,而是引擎里的一等状态:
节点类型:human_approval
执行到该节点 → 状态置 waiting_approval,执行器继续跑其他无关节点
审批 API:POST /approve /reject(带审批人、意见,落审计)
通过 → 节点 done,下游恢复调度
拒绝 → 可配置:终止整图 / 跳到补偿分支
超时策略:升级给上级 / 默认拒绝 / 默认通过(按业务风险选,默认拒绝最安全)
工程要点:
- 审批人是审计的一部分:谁批的、什么时候、基于什么材料(节点输入快照),必须不可变落库——事后追责与合规都靠它;
- 审批材料要自足:审批人可能几天后才处理,节点输入必须做成当时快照(后来被 Agent 改过的文件不能作为"当时的依据");
- 断点续跑:waiting_approval 可以持续数天,引擎重启后必须原样恢复等待状态——这要求审批状态与图状态一样落库(见 8.3)。
设计思维:HITL 节点放在不可逆动作之前(发邮件、付款、删数据、对外发布),而不是"每个 LLM 输出之后"。前者是业务风险边界,后者只是噪音审批——审批疲劳会让审批退化成无脑点确认。
8.5 事件与回放
执行过程要对三方面可观测(人、前端、调试):事件溯源是标准答案。
事件总线(EventBus)
├── 持久事件(execution_started / node_completed / approval_requested / …)
│ → 落库,带单调递增序号,支持 Last-Event-ID 续传回放
└── 瞬态事件(node_output_delta 流式 token / thinking / 进度条)
→ 只广播不落盘(量大、无回放价值),断线期间的瞬态丢失可接受
为什么瞬态不落盘:一个 50 节点的工作流可能有上百万条 delta 事件,落盘成本与价值完全不成比例。回放的语义要向前端明确:重连后恢复的是结构性状态(哪些节点 done/failed),最终结果以引擎的终态为准——前端断线期间的流式内容丢失是可接受的设计,不是 bug。
8.6 工作流 vs Agent:边界速查
| 场景 | 形态 |
|---|---|
| 步骤固定、每步简单 | 纯 workflow |
| 步骤固定、其中一步是开放任务("研究竞品并给建议") | workflow + 该节点内嵌 Agent |
| 步骤本身不可预测 | 自主 Agent(带预算与沙箱) |
| 不可逆动作 | 一律放在 workflow 节点 + human_approval,永不让自主 Agent 直接触发 |
最后一行值得画重点:自主 Agent 触发不可逆动作(发邮件/转账/删库)是生产事故的第一来源,因为它把"模型的一次概率性输出"直接连到了"现实世界的不可逆后果"。
实现作业
- 把 8.2 的最小执行器跑起来:实现一个三节点流程(读数据 → LLM 生成摘要 → 写报告),故意让中间节点抛异常,验证重试与级联 skip 的终态;
- 加持久化:节点状态与输出落 SQLite(write-ahead),进程 kill 后重启能从断点恢复(幂等节点直接重跑);
- 加
human_approval节点类型:waiting 状态 + CLI 审批命令 + 拒绝终止路径,重启进程验证审批等待能恢复; - 给执行器加事件流(print 带序号即可),体验"持久 vs 瞬态"事件的不同生命周期。
深入材料
- Temporal 文档(docs.temporal.io,生产级工作流引擎的持久化与恢复设计范本,Durable Execution 思想)
- Argo Workflows(K8s 原生 DAG 的调度实现参考)
- Anthropic, Building effective agents(workflow vs agent 谱系)
读者留言
COMMENTS 暂无还没有留言,来说第一句?