1. 为什么需要图式工作流编排多智能体
在传统的大模型应用开发中,我们常常采用线性链式调用方式:一个智能体完成任务后,将结果传递给下一个智能体。这种方式在处理简单任务时表现尚可,但当面对复杂业务场景时,就会暴露出诸多问题。
想象一下医院急诊室的运作场景。如果采用线性流程:分诊→检查→诊断→治疗→出院,这种僵化的流程根本无法应对真实医疗场景的复杂性。患者可能需要多次检查、会诊,治疗过程中可能出现并发症需要调整方案。同样,多智能体系统在处理复杂任务时,也需要类似的灵活性。
线性链式调用的主要痛点包括:
- 无法处理分支逻辑(根据条件选择不同执行路径)
- 难以实现循环迭代(如自动重试、优化循环)
- 并行处理能力有限
- 状态管理混乱
- 调试困难
而图式工作流(Graph Workflow)则像城市交通网络,智能体是各个站点,边是连接它们的道路。你可以:
- 设置单行道(顺序执行)
- 建立环线(循环逻辑)
- 构建立交桥(并行处理)
- 安排交警(条件路由)
这种架构天然适合复杂业务场景。以电商客服系统为例,一个完整的工作流可能包含:
- 意图识别
- 根据意图分支:
- 售前咨询→产品推荐
- 售后问题→订单查询→问题分类→(简单问题自动处理/复杂问题转人工)
- 满意度评价
- 根据评价决定是否转人工回访
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. LangGraph核心架构解析
LangGraph的架构设计体现了"关注点分离"的经典软件工程原则。它将工作流中的三个核心要素解耦:
2.1 节点(智能体)
节点是工作流中的处理单元,每个节点专注于完成特定任务。在LangGraph中,节点具有以下特点:
- 单一职责:每个节点只做一件事
- 无状态:节点不维护内部状态,所有状态通过参数传入
- 明确接口:输入输出类型严格定义
python复制def diagnose_agent(state: IncidentState) -> dict:
"""诊断节点实现示例"""
# 调用监控系统API获取指标
metrics = fetch_metrics(state["incident_id"])
return {"current_metrics": metrics} # 只返回增量更新
2.2 边(路由逻辑)
边定义了节点间的流转规则,分为两种类型:
- 静态边:固定不变的转移路径
python复制workflow.add_edge("diagnose", "plan_fix") - 条件边:根据状态动态决定路径
python复制def route_function(state): if state["severity"] == "high": return "emergency_procedure" return "standard_procedure" workflow.add_conditional_edges( "classify", route_function, {"emergency_procedure": emergency_node, "standard_procedure": standard_node} )
2.3 状态(共享上下文)
状态是整个工作流的"中央数据库",采用类型化设计确保数据结构一致:
python复制from typing import TypedDict, List
from datetime import datetime
class CustomerServiceState(TypedDict):
session_id: str
user_query: str
intent: str
products: List[dict]
current_step: str
created_at: datetime
retry_count: int
状态管理的关键特性:
- 不可变性:节点获得的是状态快照
- 原子更新:LangGraph负责合并增量更新
- 版本控制:每次状态变更都生成新版本
3. 实战:构建客服工单处理系统
让我们通过一个完整的实例,演示如何使用LangGraph构建生产级的智能体系统。
3.1 系统需求分析
假设我们要为电商平台构建智能工单处理系统,需求如下:
- 自动处理80%的常见问题
- 复杂问题转人工并提供处理建议
- 支持多轮对话
- 处理超时自动升级
- 记录完整处理过程
3.2 定义状态结构
首先设计状态结构,这是工作流的基础:
python复制from typing import TypedDict, List, Literal, Optional
class TicketState(TypedDict):
ticket_id: str
customer_id: str
initial_query: str
current_step: Literal[
"classify",
"handle_simple",
"prepare_complex",
"human_review",
"confirm_resolution"
]
category: Optional[Literal["return", "refund", "damage", "other"]]
solution: Optional[str]
is_resolved: bool
conversation: List[dict]
wait_time: int # 分钟
retries: int
3.3 构建工作流图
然后定义完整的工作流:
python复制from langgraph.graph import StateGraph, END
# 初始化图
workflow = StateGraph(TicketState)
# 添加节点
workflow.add_node("classify", classify_ticket)
workflow.add_node("handle_simple", handle_simple_issue)
workflow.add_node("prepare_complex", prepare_complex_case)
workflow.add_node("human_review", await_human_review)
workflow.add_node("confirm_resolution", confirm_with_customer)
# 设置边
workflow.add_edge("classify", "handle_simple")
workflow.add_edge("classify", "prepare_complex")
workflow.add_edge("prepare_complex", "human_review")
workflow.add_edge("handle_simple", "confirm_resolution")
workflow.add_edge("human_review", "confirm_resolution")
# 条件边:确认解决结果
workflow.add_conditional_edges(
"confirm_resolution",
lambda s: "resolved" if s["is_resolved"] else "retry",
{"resolved": END, "retry": "classify"}
)
# 超时检查边
workflow.add_conditional_edges(
"human_review",
lambda s: "timeout" if s["wait_time"] > 30 else "continue",
{"timeout": "escalate", "continue": "human_review"}
)
workflow.set_entry_point("classify")
3.4 实现关键节点
以分类节点为例展示具体实现:
python复制def classify_ticket(state: TicketState) -> dict:
"""工单分类节点"""
from some_llm import classify
# 调用大模型进行分类
response = classify(
prompt="分类用户问题",
query=state["initial_query"],
categories=["return", "refund", "damage", "other"]
)
# 记录对话历史
new_entry = {
"role": "system",
"content": f"分类结果: {response['category']}",
"timestamp": datetime.now().isoformat()
}
return {
"category": response["category"],
"conversation": state["conversation"] + [new_entry],
"current_step": "handle_simple" if response["is_simple"] else "prepare_complex"
}
3.5 运行与监控
编译并运行工作流:
python复制# 编译工作流
app = workflow.compile()
# 初始化工单
initial_state = {
"ticket_id": "T-1001",
"customer_id": "C-789",
"initial_query": "我收到的商品有破损,想要退货",
"current_step": "classify",
"conversation": [],
"is_resolved": False,
"wait_time": 0,
"retries": 0
}
# 执行工作流
for step in app.stream(initial_state):
print(f"当前节点: {step['current_step']}")
print(f"状态: {step.keys()}")
# 在真实系统中,这里可以添加监控逻辑
monitor_metrics(step)
4. 高级特性与最佳实践
4.1 并行执行优化
对于可以并行处理的任务,LangGraph提供了优雅的实现方式:
python复制# 添加并行节点
workflow.add_node("check_inventory", check_inventory)
workflow.add_node("check_policy", check_return_policy)
workflow.add_node("check_customer", check_customer_history)
# 从classify同时指向三个节点
workflow.add_edge("classify", "check_inventory")
workflow.add_edge("classify", "check_policy")
workflow.add_edge("classify", "check_customer")
# 添加聚合节点
workflow.add_node("evaluate_case", evaluate_case)
workflow.add_edge("check_inventory", "evaluate_case")
workflow.add_edge("check_policy", "evaluate_case")
workflow.add_edge("check_customer", "evaluate_case")
4.2 状态Reducer模式
当多个节点需要更新同一状态字段时,使用Reducer保持一致性:
python复制from typing import Annotated
from operator import add
class EvaluationState(TypedDict):
inventory_status: str
policy_decision: str
customer_rating: int
# 使用Reducer合并多个评分
risk_scores: Annotated[list[float], add]
def check_inventory(state) -> dict:
# 检查库存逻辑...
return {"risk_scores": [0.2]} # 返回增量
def check_policy(state) -> dict:
# 检查政策逻辑...
return {"risk_scores": [0.3]} # 返回增量
4.3 检查点与恢复
LangGraph自动创建检查点,支持故障恢复:
python复制# 从检查点恢复
checkpoint_id = get_last_checkpoint(ticket_id="T-1001")
recovered = app.recover(checkpoint_id)
# 继续执行
for step in recovered.stream():
process_step(step)
4.4 超时与熔断机制
在生产环境中必须添加保护措施:
python复制from datetime import datetime, timedelta
def with_timeout(agent_fn, timeout=30):
"""为节点添加超时保护的装饰器"""
def wrapped(state):
start = datetime.now()
result = agent_fn(state)
duration = (datetime.now() - start).seconds
if duration > timeout:
return {"timed_out": True, "current_step": "fallback"}
return result
return wrapped
# 应用超时保护
workflow.add_node("llm_call", with_timeout(llm_agent))
5. 生产环境部署建议
5.1 性能优化技巧
-
节点级缓存:
python复制from functools import lru_cache @lru_cache(maxsize=1000) def classify_ticket(query: str) -> dict: # 缓存相同查询的结果 pass -
批量处理:
python复制def batch_process_tickets(tickets: list[TicketState]): # 合并相似工单批量处理 pass -
异步执行:
python复制async def async_agent(state): # 使用异步IO提高吞吐量 pass
5.2 监控与可观测性
建议监控以下指标:
- 节点执行时间
- 工作流完成率
- 循环次数分布
- 分支路径分布
- 错误率
使用Prometheus示例:
python复制from prometheus_client import Summary
NODE_TIME = Summary('node_processing_time', 'Time spent processing nodes')
@NODE_TIME.time()
def handle_simple_issue(state):
# 节点逻辑
pass
5.3 版本控制策略
工作流也需要版本管理:
- 使用Git管理图定义
- 状态结构变更时考虑兼容性
- 部署新版本前进行影子测试
- 保留旧版本一定时间以便回滚
python复制# 版本化工作流定义
class TicketWorkflowV1:
# 第一版实现
pass
class TicketWorkflowV2(TicketWorkflowV1):
# 向后兼容的修改
pass
6. 与其他方案的对比
6.1 对比LangChain
LangChain更适合线性流程,而LangGraph优势在于:
- 复杂流程可视化
- 内置状态管理
- 更好的调试支持
- 更灵活的路由
6.2 对比Airflow
虽然都是工作流引擎,但专注点不同:
- LangGraph专为AI智能体优化
- 内置LLM相关功能
- 状态管理更轻量
- 更适合实时交互场景
6.3 对比自定义实现
自行开发类似系统需要考虑:
- 状态持久化
- 错误恢复
- 并发控制
- 监控集成
使用LangGraph可以节省约70%的相关开发工作量
