LangGraph:大模型工作流编排框架解析与应用
1. LangGraph 在大模型架构中的定位
LangGraph 是 LangChain 生态系统中专门用于复杂工作流编排的框架。与 LangChain 提供的链式调用不同,LangGraph 引入了基于图的执行模型,更适合处理需要条件分支、循环和并行执行的 AI 应用场景。这种架构特别适合当前大模型应用中常见的多步骤推理、多智能体协作等需求。
在技术实现上,LangGraph 将工作流抽象为有向图(Directed Graph),其中节点代表处理单元(可以是 LLM 调用、工具使用或自定义函数),边代表执行路径。这种设计带来了几个关键优势:
- 动态路由:可以根据前序节点的输出结果动态选择后续路径
- 状态管理:通过全局的"状态对象"在不同节点间传递和修改数据
- 循环支持:内置机制可以实现类似 while 循环的重复执行
- 并行执行:多个不依赖的节点可以同时运行
提示:虽然 LangGraph 和 LangChain 都来自同一生态,但 LangGraph 不是 LangChain 的替代品,而是专门解决工作流编排这个特定问题的补充方案。
2. 核心概念与架构解析
2.1 图结构定义
LangGraph 的核心抽象是StateGraph,构建一个工作流通常包含以下步骤:
from langgraph.graph import StateGraph # 定义状态结构 from typing import TypedDict, List class GraphState(TypedDict): input: str intermediate_results: List[str] final_output: str # 创建图实例 workflow = StateGraph(GraphState)状态类型使用 Python 的 TypedDict 定义,这为整个工作流提供了类型安全的上下文对象。每个节点都可以读取和修改这个状态对象的任意字段。
2.2 节点与边
节点是工作流的基本执行单元,通常包装了 LLM 调用、工具使用或业务逻辑:
def retrieve_node(state: GraphState): # 检索增强生成(RAG)的检索步骤 documents = vector_store.similarity_search(state["input"]) return {"intermediate_results": [doc.page_content for doc in documents]} def generate_node(state: GraphState): # 生成步骤 prompt = ChatPromptTemplate.from_template("基于以下内容回答问题:{context}\n问题:{question}") chain = prompt | llm response = chain.invoke({ "context": "\n".join(state["intermediate_results"]), "question": state["input"] }) return {"final_output": response.content} # 添加节点 workflow.add_node("retriever", retrieve_node) workflow.add_node("generator", generate_node)边的定义决定了工作流的走向:
# 设置入口点 workflow.set_entry_point("retriever") # 定义节点连接 workflow.add_edge("retriever", "generator") workflow.add_edge("generator", END) # END是特殊终止节点 # 编译可执行图 app = workflow.compile()2.3 条件路由
更复杂的场景需要条件分支,LangGraph 通过add_conditional_edges实现:
def should_continue(state: GraphState): # 根据生成内容决定是否继续 if "需要更多信息" in state["final_output"]: return "expand_query" else: return END workflow.add_conditional_edges( "generator", should_continue, {"expand_query": "query_expander", END: END} )3. 多智能体协作实现
LangGraph 特别适合实现多智能体(Multi-Agent)系统。下面是一个评审-修订模式的实现示例:
3.1 智能体定义
class AgentState(TypedDict): draft: str feedback: List[str] final_version: str def writer_agent(state: AgentState): prompt = """根据以下反馈改进文档: 反馈:{feedback} 当前版本:{draft} 输出修订后的版本:""" response = llm.invoke(prompt.format( feedback="\n".join(state["feedback"]), draft=state["draft"] )) return {"draft": response.content} def reviewer_agent(state: AgentState): prompt = """评审以下文档并提出具体改进建议: 文档:{draft}""" response = llm.invoke(prompt.format(draft=state["draft"])) return {"feedback": [response.content]}3.2 协作工作流
workflow = StateGraph(AgentState) workflow.add_node("write", writer_agent) workflow.add_node("review", reviewer_agent) workflow.set_entry_point("write") # 第一轮:写 → 审 workflow.add_edge("write", "review") # 条件路由:最多三轮迭代 def should_iterate(state: AgentState): if len(state["feedback"]) < 3: return "write" else: return "finalize" workflow.add_conditional_edges( "review", should_iterate, {"write": "write", "finalize": "finalize"} ) def finalize(state: AgentState): return {"final_version": state["draft"]} workflow.add_node("finalize", finalize) workflow.add_edge("finalize", END)4. 高级特性与调试技巧
4.1 持久化与检查点
LangGraph 支持工作流状态的持久化,这对长时间运行的流程特别有用:
from langgraph.checkpoint import MemorySaver app = workflow.compile( checkpointer=MemorySaver(), interrupt_before=["review"] # 在评审前允许中断 ) # 运行并保存 config = {"configurable": {"thread_id": "user123"}} result1 = app.invoke({"draft": "初稿"}, config) # 后续恢复执行 result2 = app.invoke(None, config) # 传入None表示继续上次状态4.2 可视化调试
LangGraph 集成了 LangSmith 提供可视化跟踪:
import os os.environ["LANGCHAIN_TRACING_V2"] = "true" os.environ["LANGCHAIN_PROJECT"] = "langgraph-tutorial" # 现在执行会生成可视化轨迹 app.invoke({"input": "我的问题"})在 LangSmith 控制台可以看到完整的执行流程图,包括每个节点的输入输出和耗时。
4.3 性能优化技巧
- 并行执行:使用
add_edge的parallel参数让不依赖的节点同时运行 - 批处理:在节点内部实现批处理逻辑减少 LLM 调用次数
- 缓存:为 LLM 节点配置缓存(如使用
langchain.cache) - 超时控制:为每个节点设置合理的超时时间
from langgraph.graph import RETRY, TIMEOUT workflow.add_node( "api_call", api_node.with_retry(max_attempts=3).with_timeout(30) )5. 生产环境最佳实践
5.1 错误处理策略
建议实现以下错误处理机制:
def safe_node(state: GraphState): try: # 正常节点逻辑 return {"key": value} except Exception as e: # 错误处理逻辑 return {"error": str(e), "retry": True} workflow.add_node("safe_step", safe_node) # 全局错误处理 def handle_error(state: GraphState): if any("error" in step for step in state.values()): return "error_handler" return "next_step" workflow.add_conditional_edges( "safe_step", handle_error, {"error_handler": "error_node", "next_step": "normal_node"} )5.2 监控指标
关键监控指标建议:
| 指标类别 | 具体指标 | 采集方式 |
|---|---|---|
| 性能 | 节点执行时间、吞吐量 | LangSmith 日志 |
| 质量 | 输出合规率、用户满意度 | 人工审核+反馈系统 |
| 可靠性 | 错误率、重试次数 | 节点错误捕获 |
| 成本 | LLM token 使用量 | API 调用日志 |
5.3 版本控制策略
对于工作流变更,推荐采用:
- 蓝绿部署:同时维护新旧版本的工作流
- 流量分流:通过配置将部分请求导向新版本
- 自动化回滚:基于监控指标自动回退有问题的版本
# 版本路由示例 def route_by_version(state: GraphState): if state.get("api_version") == "v2": return "v2_workflow" else: return "v1_workflow"在实际项目中,我们通常会将这些工作流定义存储在版本控制的配置文件中,与 CI/CD 管道集成。
