Agent 工程 · 第 8 章|工作流引擎与 HITL:DAG 执行器、审批、事件回放

第 8 章 · 工作流引擎与 Human-in-the-Loop

Agent 的自主性是资产也是负债。工作流(DAG 引擎)是把确定性控制流从模型手里拿回来的基础设施;Human-in-the-Loop(HITL)是在关键决策点把控制权交还给人。两者合起来回答一个问题:自主性放到什么程度,在哪里设闸门。

8.1 为什么需要工作流引擎

考虑一个生产任务:"分析这份财报 → 生成摘要 → 风控审查 → 发邮件"。三种实现:

  1. 全自主 Agent:让模型自己决定调工具的顺序。风险:顺序不可控(没审查就发了邮件)、中途失败无法恢复、成本不可预测;
  2. 硬编码脚本:代码里写死四步。可靠,但每一步内部的 LLM 调用细节、重试、并行都要手写,改一个步骤 = 发一个版本;
  3. 工作流引擎:步骤与依赖声明式定义(图),执行、重试、状态、恢复、事件全部由引擎统一提供。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 节点——无法确认是否已执行(可能执行到一半崩了)。两种处理:

  1. 节点声明幂等(同样输入重复执行结果一致)→ 直接重跑;
  2. 不幂等(发邮件、付款)→ 执行器需要"执行前意向日志"(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 触发不可逆动作(发邮件/转账/删库)是生产事故的第一来源,因为它把"模型的一次概率性输出"直接连到了"现实世界的不可逆后果"。

实现作业

  1. 把 8.2 的最小执行器跑起来:实现一个三节点流程(读数据 → LLM 生成摘要 → 写报告),故意让中间节点抛异常,验证重试与级联 skip 的终态;
  2. 加持久化:节点状态与输出落 SQLite(write-ahead),进程 kill 后重启能从断点恢复(幂等节点直接重跑);
  3. 加 human_approval 节点类型:waiting 状态 + CLI 审批命令 + 拒绝终止路径,重启进程验证审批等待能恢复;
  4. 给执行器加事件流(print 带序号即可),体验"持久 vs 瞬态"事件的不同生命周期。

深入材料

  • Temporal 文档(docs.temporal.io,生产级工作流引擎的持久化与恢复设计范本,Durable Execution 思想)
  • Argo Workflows(K8s 原生 DAG 的调度实现参考)
  • Anthropic, Building effective agents(workflow vs agent 谱系)
← 返回资讯列表

读者留言

COMMENTS 暂无
仅本站原创文章开放留言 · 请勿留下手机号、邮箱等个人信息

还没有留言,来说第一句?