Agent 工作流编排设计:DAG、状态机、暂停恢复与长期运行任务

深入探索 LLM Agent 的工作流编排模型:DAG(有向无环图)编排、状态机驱动、事件触发与暂停恢复。 涵盖 LangGraph 的图式编排、Temporal 的长任务持久化、以及事件总线驱动的异步 Agent 架构。 附完整的 LangGraph + Temporal 集成案例,含可运行代码与可视化。

简单的 ReAct 循环可以处理单次任务,但在生产环境中,Agent 往往需要处理跨越小时甚至数天的复杂流程。本文系统拆解工作流编排的三种核心模型,以及如何用现代工作流引擎实现稳健的长期运行任务。


1. 三种编排模型

1.1 DAG(有向无环图)

         用研 → 设计
           ↘      ↓
             → 开发 → 测试
                      ↓
                   上线
  • 优点:执行顺序清晰、并行度可控、可视化友好
  • 缺点:不支持循环(需要补充状态机)

1.2 状态机(State Machine)

[待处理] → [处理中] → [待审核]
              ↓          ↓
          [已完成]    [已退回]
  • 优点:明确的内部状态 + 转换条件、易于审计
  • 缺点:手工定义所有状态成本高

1.3 事件驱动(Event-Driven)

事件总线
├── Agent-A 监听 "user.created" → 发送欢迎邮件
├── Agent-B 监听 "order.placed" → 检查库存
└── Agent-C 监听 "payment.failed" → 安排人工介入
  • 优点:解耦极强、高扩展、适合微服务环境
  • 缺点:调试困难、事序难以保证

2. LangGraph 深度实战

2.1 基本图结构

from langgraph.graph import StateGraph, END
from typing import TypedDict, Any
import json

# 共享状态
class WorkflowState(TypedDict):
    input_text: str
    intent: str
    extracted_entities: dict
    api_response: Any
    final_answer: str

# 节点函数
def intent_recognition(state: WorkflowState):
    # 意图识别 Agent
    return {"intent": "book_flight"}

def entity_extraction(state: WorkflowState):
    # 实体提取
    return {"extracted_entities": {"origin": "上海", "destination": "北京"}}

def api_call(state: WorkflowState):
    # API 查询
    return {"api_response": {"flights": ["MU5161", "CA1850"]}}

def generate_response(state: WorkflowState):
    # 回复生成
    return {"final_answer": "为您找到 2 班航班..."}

# 构建工作流
builder = StateGraph(WorkflowState)
builder.add_node("intent", intent_recognition)
builder.add_node("extract", entity_extraction)
builder.add_node("api", api_call)
builder.add_node("respond", generate_response)

# 定义边(顺序执行)
builder.add_edge("intent", "extract")
builder.add_edge("extract", "api")
builder.add_edge("api", "respond")
builder.add_edge("respond", END)
builder.set_entry_point("intent")

# 编译执行
workflow = builder.compile()
result = workflow.invoke({"input_text": "帮我订上海到北京的机票"})

2.2 条件分支

def route_by_intent(state):
    """根据意图动态路由"""
    return state["intent"]  # 返回字符串决定下一个节点

builder.add_conditional_edges(
    "intent",
    route_by_intent,
    {
        "book_flight": "extract",
        "check_weather": "api",
        "general_chat": "respond",
    }
)

2.3 并行节点

from langgraph.graph import ANNOTATED_STATE
import operator

class ParallelState(TypedDict):
    queries: list[str]
    results: Annotated[list, operator.add]  # 并行结果累加

# researchA 和 researchB 同时运行
builder.add_node("research_a", research_agent_a)
builder.add_node("research_b", research_agent_b)
builder.add_node("synthesize", synthesis_agent)

builder.add_edge("synthesize", "synthesize")

2.4 循环(回头边)

def check_quality(state):
    if state["quality_score"] < 0.8:
        return "retry"  # → 回到 edit 节点
    return "done"      # → END

builder.add_conditional_edges(
    "evaluate",
    check_quality,
    {"retry": "edit", "done": END}
)

3. 暂停与恢复机制

长期任务需要在中间步骤暂停等待外部输入,并在收到信号后恢复。

3.1 LangGraph Interrupt

from langgraph.types import interrupt

async def human_review(state):
    # 暂停等待人类审核
    result = interrupt(
        value={"draft": state["draft"]},
        # 弹出 UI 让用户确认
    )
    # result 是人类返回的数据
    return {"approved": result["approved"]}

builder.add_node("human_review", human_review)
builder.add_edge("draft", "human_review")
builder.add_edge("human_review", "publish")

3.2 基于数据库的持久化恢复

import uuid
from datetime import datetime

class PersistentWorkflow:
    def __init__(self, db):
        self.db = db

    async def start(self, workflow_def, initial_state):
        run_id = str(uuid.uuid4())
        await self.db.execute(
            "INSERT INTO workflow_runs (id, status, state) VALUES (?, 'running', ?)",
            (run_id, json.dumps(initial_state))
        )
        return run_id

    async def checkpoint(self, run_id, state):
        await self.db.execute(
            "UPDATE workflow_runs SET state = ?, checkpoint_time = ? WHERE id = ?",
            (json.dumps(state), datetime.utcnow().isoformat(), run_id)
        )

    async def resume(self, run_id, human_input):
        row = await self.db.fetchone(
            "SELECT * FROM workflow_runs WHERE id = ?", (run_id,)
        )
        state = json.loads(row["state"])
        # 注入人类输入并继续执行
        state["human_input"] = human_input
        return await self._run_from_checkpoint(state)

4. Temporal:生产级长任务编排

Temporal 是 Uber 开源的工作流引擎,特别适合跨小时、跨天的长期任务。

4.1 为什么用 Temporal?

能力LangGraphTemporal
长任务耐久手动持久化自动 Checkpoint
分布式执行单次运行多 Worker 并行
定时触发不支持Cron 内置
重试策略手动声明式
可视化 UI基础完整工作流 UI
学习曲线平缓中等

4.2 Temporal + LLM Agent 集成示例

from temporalio import workflow
from temporalio.client import Client
import asyncio

# Workflow 定义
@workflow.defn
class ResearchAgentWorkflow:
    @workflow.run
    async def run(self, research_topic: str) -> str:
        # Step 1: 搜索(可中断)
        search_results = await workflow.execute_activity(
            "search_activity",
            research_topic,
            start_to_close_timeout=timedelta(minutes=1),
            retry_policy=RetryPolicy(maximum_attempts=3),
        )

        # Step 2: 分析(可并行)
        analyses = await asyncio.gather(*[
            workflow.execute_activity("analysis_activity", result)
            for result in search_results
        ])

        # Step 3: 人类审核(长时间等待)
        draft = await workflow.execute_activity("synthesis_activity", analyses)
        
        # 暂停等待人类评审
        review = await workflow.execute_activity(
            "request_human_review",
            draft,
            start_to_close_timeout=timedelta(days=7),  # 人类有 7 天时间
            schedule_to_close_timeout=timedelta(days=30),
        )

        # Step 4: 修改
        if not review["approved"]:
            draft = await workflow.execute_activity(
                "revise_activity",
                {"draft": draft, "feedback": review["feedback"]},
            )

        # Step 5: 发布
        return await workflow.execute_activity("publish_activity", draft)

4.3 Temporal Worker

from temporalio.worker import Worker

async def main():
    client = await Client.connect("temporal:7233")
    
    activities = [search_activity, analysis_activity, synthesis_activity]
    
    worker = Worker(
        client,
        task_queue="agent-research-queue",
        workflows=[ResearchAgentWorkflow],
        activities=activities,
    )
    await worker.run()

5. 事件总线驱动的异步 Agent

from dataclasses import dataclass
from typing import Callable, List
import asyncio

@dataclass
class AgentEvent:
    event_type: str  # "order.placed", "review.pending", etc.
    payload: dict
    correlation_id: str

class EventBus:
    def __init__(self):
        self.subscribers: dict[str, List[Callable]] = {}
        self.event_queue = asyncio.Queue()

    def subscribe(self, event_type: str, handler: Callable):
        self.subscribers.setdefault(event_type, []).append(handler)

    async def emit(self, event: AgentEvent):
        await self.event_queue.put(event)

    async def run(self):
        while True:
            event = await self.event_queue.get()
            handlers = self.subscribers.get(event.event_type, [])
            for handler in handlers:
                asyncio.create_task(handler(event))

# Agent 订阅事件
bus = EventBus()

def order_agent(event: AgentEvent):
    print(f"处理订单: {event.payload}")
    # 发起支付、检查库存...

bus.subscribe("order.placed", order_agent)

# 触发事件
asyncio.run(bus.emit(AgentEvent("order.placed", {"sku": "A001", "qty": 2}, "uuid-123")))

6. 架构选型总结

任务持续时间?
├── <10 秒 → LangChain + 同步调用
│
├── <10 分钟 → LangGraph(图编排)
│   ├── 需要循环? → 状态机 + LangGraph
│   └── 人类介入? → LangGraph Interrupt
│
├── <24 小时 → LangGraph + 持久化 Checkpoint + 重试
│
└── >24 小时 → Temporal(分布式 + 自动耐久)

FAQ

Q: 什么时候用 LangGraph,什么时候用 Temporal?
A: 如果你的 Agent 任务能在 10 分钟内完成且不需要外部长期等待(如人类审批),LangGraph 更轻量。如果需要跨越数小时的任务、Cron 定时触发、或跨团队共享工作流状态,选 Temporal。

Q: Agent 工作流和 RPA(机器人流程自动化)有什么区别?
A: RPA 是确定性流程(步骤 100% 预设),Agent 工作流是意图驱动的动态流程(模型自主决定下一步)。Agent 更灵活但不可预测性也更高;在需要严格合规审计的场景(如金融审批),RPA 仍不可替代。

Q: 人类介入(HITL)的设计要点?
A: 三个原则:1)设置明确的 Timeout(如 48h 无响应则走降级策略);2)提供上下文摘要(不要让人类阅读几百条消息);3)保留审批链的不可篡改日志(用于合规审计)。

📂 相关专题:

下一篇 →

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「llm」更多文章