简单的 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?
| 能力 | LangGraph | Temporal |
|---|---|---|
| 长任务耐久 | 手动持久化 | 自动 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)保留审批链的不可篡改日志(用于合规审计)。
📂 相关专题:
- AI 智能体架构设计 — 五大核心组件详解
- 多智能体协作与编排 — CrewAI、AutoGen、LangGraph 实战对比
- LLM Agent 与人协作模式 — HITL 最佳实践
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。