LangGraph 从入门到精通:AI Agent 编排框架深度解析与实战案例
·
LangGraph 从入门到精通:AI Agent 编排框架深度解析与实战案例
摘要:本文系统讲解 LangGraph 框架的核心概念、架构设计原理,并通过智能客服工单系统的完整实战案例,展示如何从传统的线性 Chain 模式演进到基于图的灵活编排。适合 AI 开发者从入门到进阶全面掌握。

目录
一、为什么需要 LangGraph?
1.1 传统 AI 应用开发的痛点
在 LangGraph 出现之前,开发者使用 LangChain 构建 AI 应用时,主要面临以下挑战:
- 线性流程限制:传统的 Chain 模式是线性的,难以处理复杂的多分支逻辑
- 状态管理混乱:缺乏统一的状态管理机制,依赖全局变量导致代码难以维护
- 流程不可视:复杂的业务逻辑难以直观展示和调试
- 缺乏中断能力:无法在流程执行过程中进行人工干预
- 错误追踪困难:节点级别的错误难以精确定位和追踪
1.2 LangGraph 的解决方案
LangGraph 基于**有向图(Directed Graph)**的设计理念,将 AI Agent 的工作流程抽象为图结构,提供了:
- ✅ 灵活的流程控制:支持条件分支、循环、并行等复杂流程
- ✅ 显式状态管理:通过 State 对象统一管理数据流转
- ✅ 可视化能力:自动生成 Mermaid 流程图
- ✅ 中断与恢复:支持在任意节点暂停并等待人工输入
- ✅ 可调试性:节点级别的错误追踪和日志记录

二、LangGraph 核心概念解析
2.1 State(状态)- 数据流转的载体
State 是 LangGraph 的核心数据结构,用于在节点之间传递和共享数据。
2.1.1 状态定义方式
from typing import TypedDict, List, Annotated
import operator
class TicketAgentState(TypedDict):
"""智能客服工单系统状态定义"""
user_id: str # 用户 ID
query: str # 用户问题
category: str # 问题分类
solution: str # 解决方案
action_approved: bool # 操作是否批准
risk_level: int # 风险等级 (1-3)
audit_log: List[str] # 审计日志
progress: str # 当前进度
error: str # 错误信息
session_id: str # 会话 ID
current_step: str # 当前步骤
# 使用 Annotated 定义状态更新策略
messages: Annotated[List[str], operator.add] # 消息列表,使用追加策略
2.1.2 状态更新策略
LangGraph 支持多种状态更新策略:
- 覆盖(overwrite):默认策略,新值直接覆盖旧值
- 追加(append):使用
operator.add,新值追加到列表 - 自定义策略:通过
Annotated定义自定义 reducer 函数
from typing import Annotated
import operator
# 示例:自定义状态合并策略
def merge_states(old: dict, new: dict) -> dict:
"""合并两个状态字典"""
return {**old, **new, 'updated_at': new.get('updated_at')}
class CustomState(TypedDict):
data: Annotated[dict, merge_states]
2.2 Node(节点)- 执行单元
Node 是图的基本执行单元,负责处理具体的业务逻辑。
2.2.1 节点类型
# 1. LLM 节点 - 调用大语言模型
async def classify_query_node(state: TicketAgentState) -> dict:
"""问题分类节点"""
from langchain_openai import ChatOpenAI
llm = ChatOpenAI(model="gpt-4o", temperature=0.7)
prompt = f"""
请将以下用户问题分类到合适的类别:
用户问题:{state['query']}
可选类别:
- 账户问题
- 支付问题
- 技术支持
- 投诉建议
- 其他
请直接返回分类结果(仅返回类别名称):
"""
response = await llm.ainvoke(prompt)
return {"category": response.content.strip()}
# 2. 工具节点 - 调用外部工具
async def search_knowledge_base(state: TicketAgentState) -> dict:
"""知识库查询节点"""
# 模拟知识库查询
solutions = {
"账户问题": "请检查您的登录状态,尝试重新登录或重置密码",
"支付问题": "请确认支付方式是否有效,联系银行核实交易状态",
"技术支持": "请提供详细错误信息和技术环境描述",
}
solution = solutions.get(state['category'], "请联系人工客服")
return {"solution": solution}
# 3. 规则节点 - 执行业务规则
async def evaluate_risk_node(state: TicketAgentState) -> dict:
"""风险等级评估节点"""
# 基于规则的风险分析
risk_keywords = ["投诉", "法律", "监管", "媒体"]
query = state['query'].lower()
risk_score = 1 # 默认低风险
for keyword in risk_keywords:
if keyword in query:
risk_score += 1
break
risk_score = min(risk_score, 3) # 限制在 1-3 之间
return {"risk_level": risk_score}
# 4. 人工节点 - 等待人工干预
async def human_review_node(state: TicketAgentState) -> dict:
"""人工审核节点"""
print(f"\n⚠️ 需要人工审核 - 会话 ID: {state['session_id']}")
print(f"用户问题:{state['query']}")
print(f"建议方案:{state['solution']}")
print(f"风险等级:{state['risk_level']}")
# 等待人工输入
approval = input("是否批准执行?(yes/no): ").strip().lower()
is_approved = approval in ['yes', 'y', '是']
return {"action_approved": is_approved}
# 5. API 节点 - 执行外部调用
async def execute_api_node(state: TicketAgentState) -> dict:
"""执行 API 调用节点"""
import asyncio
# 模拟 API 调用
await asyncio.sleep(1)
result = f"✅ 已执行操作:{state['solution']}"
audit_entry = f"[{state['session_id']}] 执行解决方案:{state['solution']}"
return {
"progress": "completed",
"audit_log": [audit_entry]
}
2.3 Edge(边)- 流程控制
Edge 定义了节点之间的转换关系,支持无条件边和条件边。
2.3.1 无条件边
# 简单的顺序连接
graph.add_edge("start_node", "process_node")
# 执行 start_node 后,直接进入 process_node
2.3.2 条件边
# 基于状态字段的条件路由
def route_based_on_risk_level(state: TicketAgentState) -> str:
"""根据风险等级路由"""
if state['risk_level'] >= 3:
return "human_review" # 高风险,需要人工审核
else:
return "execute_api" # 低风险,直接执行
def route_based_on_review_approval(state: TicketAgentState) -> str:
"""根据审核结果路由"""
if state['action_approved']:
return "execute_api" # 审核通过,执行操作
else:
return "end_session" # 审核拒绝,结束会话
# 添加条件边
graph.add_conditional_edges(
"evaluate_risk",
route_based_on_risk_level,
{
"human_review": "human_review",
"execute_api": "execute_api"
}
)
graph.add_conditional_edges(
"human_review",
route_based_on_review_approval,
{
"execute_api": "execute_api",
"end_session": "end_session"
}
)
2.4 Graph(图)- 编排引擎
Graph 是 LangGraph 的编排核心,负责组合节点和边,形成完整的工作流。
2.4.1 图的构建流程
from langgraph.graph import StateGraph
# 1. 创建图实例
graph = StateGraph(TicketAgentState)
# 2. 添加节点
graph.add_node("start", start_node)
graph.add_node("classify", classify_query_node)
graph.add_node("generate_solution", generate_solution_node)
graph.add_node("evaluate_risk", evaluate_risk_node)
graph.add_node("human_review", human_review_node)
graph.add_node("execute_api", execute_api_node)
graph.add_node("audit_log", audit_log_node)
graph.add_node("end", end_node)
# 3. 设置入口点
graph.set_entry_point("start")
# 4. 添加边(无条件)
graph.add_edge("start", "classify")
graph.add_edge("classify", "generate_solution")
graph.add_edge("generate_solution", "evaluate_risk")
# 5. 添加条件边
graph.add_conditional_edges(
"evaluate_risk",
route_based_on_risk_level,
{
"human_review": "human_review",
"execute_api": "execute_api"
}
)
graph.add_conditional_edges(
"human_review",
route_based_on_review_approval,
{
"execute_api": "execute_api",
"end_session": "end_session"
}
)
graph.add_edge("execute_api", "audit_log")
graph.add_edge("audit_log", "end")
# 6. 设置结束点
graph.add_edge("end_session", "end")
# 7. 编译图
app = graph.compile()
2.4.2 图的执行
# 同步执行
initial_state = {
"user_id": "user_12345",
"query": "我的账户无法登录,这是投诉!",
"session_id": "session_789",
"messages": []
}
final_state = app.invoke(initial_state)
print(f"最终状态:{final_state}")
# 异步执行
final_state = await app.ainvoke(initial_state)
# 流式执行(逐步获取中间状态)
for chunk in app.stream(initial_state):
for node_name, node_output in chunk.items():
print(f"节点 {node_name} 输出:{node_output}")
三、与传统框架的对比优势
3.1 可视化能力
LangGraph 支持将图结构导出为 Mermaid 格式,便于文档化和分享。
# 导出为 Mermaid 图表
mermaid_code = app.get_graph().draw_mermaid()
print(mermaid_code)

3.2 中断与恢复(Human-in-the-Loop)
from langgraph.graph import MessagesState
# 配置中断点
app = graph.compile()
# 在特定节点前中断
config = {"configurable": {"thread_id": "thread_123"}}
app.update_state(config, {"progress": "paused"})
# 恢复执行
resumed_state = app.invoke(None, config)
3.3 显式状态管理
# 传统方式:依赖全局变量
global_state = {}
def node_a():
global_state['data'] = process()
return global_state['data']
# LangGraph 方式:显式状态传递
def node_a(state: MyState) -> dict:
new_data = process()
return {"data": new_data}
3.4 调试与追踪
from langsmith import Client
# 集成 LangSmith 进行追踪
client = Client()
# 记录执行轨迹
run_id = client.create_run(
name="ticket_agent_execution",
run_type="chain",
inputs=initial_state
)
try:
final_state = app.invoke(initial_state)
client.update_run(run_id, outputs=final_state, status="success")
except Exception as e:
client.update_run(run_id, error=str(e), status="error")
raise
四、实战案例:智能客服工单系统
4.1 系统架构设计
本案例实现一个完整的智能客服工单系统,具备以下功能:
- 自动分类:识别用户问题类型
- 方案生成:根据分类生成解决方案
- 风险评估:评估操作风险等级
- 人工审核:高风险操作需要人工确认
- 执行操作:调用 API 执行解决方案
- 审计日志:记录完整操作轨迹

4.2 完整代码实现
4.2.1 状态定义
from typing import TypedDict, List, Annotated
import operator
import uuid
from datetime import datetime
class TicketAgentState(TypedDict):
"""智能客服工单系统状态"""
user_id: str
query: str
category: str
solution: str
action_approved: bool
risk_level: int
audit_log: List[str]
progress: str
error: str
session_id: str
current_step: str
timestamp: Annotated[str, operator.add] # 时间戳列表
# 使用 Pydantic 替代方案
# from pydantic import BaseModel
# class TicketAgentState(BaseModel):
# user_id: str
# query: str
# ...
4.2.2 节点实现
import asyncio
from langchain_openai import ChatOpenAI
# 1. 开始节点
async def start_node(state: TicketAgentState) -> dict:
"""初始化会话"""
session_id = f"session_{uuid.uuid4().hex[:8]}"
timestamp = datetime.now().isoformat()
return {
"session_id": session_id,
"progress": "started",
"current_step": "initialization",
"timestamp": [f"会话开始:{timestamp}"]
}
# 2. 问题分类节点
async def classify_query_node(state: TicketAgentState) -> dict:
"""使用 LLM 分类用户问题"""
llm = ChatOpenAI(model="gpt-4o", temperature=0.3)
prompt = f"""
请将以下用户问题分类到合适的类别:
用户问题:{state['query']}
可选类别:
1. 账户问题 - 登录、注册、密码等
2. 支付问题 - 付款、退款、账单等
3. 技术支持 - 功能使用、错误排查等
4. 投诉建议 - 意见反馈、投诉等
5. 其他 - 不属于以上类别
请只返回类别名称,不要其他内容。
"""
response = await llm.ainvoke(prompt)
category = response.content.strip()
timestamp = datetime.now().isoformat()
audit_entry = f"[{state.get('session_id', 'unknown')}] 问题分类:{category}"
return {
"category": category,
"current_step": "classified",
"audit_log": [audit_entry],
"timestamp": [f"分类完成:{timestamp}"]
}
# 3. 生成解决方案节点
async def generate_solution_node(state: TicketAgentState) -> dict:
"""根据分类生成解决方案"""
# 方案知识库
solution_db = {
"账户问题": """
账户问题解决方案:
1. 检查网络连接
2. 尝试重新登录
3. 如忘记密码,使用"忘记密码"功能重置
4. 如账号被锁定,联系管理员解锁
""",
"支付问题": """
支付问题解决方案:
1. 确认支付方式余额充足
2. 检查银行卡是否过期
3. 联系银行确认交易状态
4. 如重复扣款,申请退款
""",
"技术支持": """
技术支持解决方案:
1. 查看官方文档和 FAQ
2. 检查系统状态和公告
3. 提供详细错误信息和复现步骤
4. 尝试清除缓存和重新安装
""",
"投诉建议": """
投诉建议处理流程:
1. 记录详细投诉内容
2. 转交客服专员处理
3. 24 小时内给予回复
4. 持续跟进直至解决
""",
}
category = state['category']
solution = solution_db.get(category, "请联系人工客服获取帮助")
timestamp = datetime.now().isoformat()
audit_entry = f"[{state.get('session_id', 'unknown')}] 生成解决方案:{category}"
return {
"solution": solution,
"current_step": "solution_generated",
"audit_log": [audit_entry],
"timestamp": [f"方案生成:{timestamp}"]
}
# 4. 风险评估节点
async def evaluate_risk_node(state: TicketAgentState) -> dict:
"""评估操作风险等级"""
query = state['query'].lower()
category = state['category']
# 风险评分规则
risk_score = 1 # 默认低风险
# 高风险关键词
high_risk_keywords = ["投诉", "法律", "监管", "媒体", "曝光", "起诉"]
medium_risk_keywords = ["退款", "赔偿", "注销", "删除数据"]
for keyword in high_risk_keywords:
if keyword in query:
risk_score = 3
break
if risk_score < 3:
for keyword in medium_risk_keywords:
if keyword in query:
risk_score = 2
break
# 特定分类自动提升风险
if category == "投诉建议":
risk_score = max(risk_score, 2)
risk_score = min(risk_score, 3) # 限制在 1-3
timestamp = datetime.now().isoformat()
audit_entry = f"[{state.get('session_id', 'unknown')}] 风险评估:等级{risk_score}"
return {
"risk_level": risk_score,
"current_step": "risk_evaluated",
"audit_log": [audit_entry],
"timestamp": [f"风险评估:{timestamp}"]
}
# 5. 人工审核节点
async def human_review_node(state: TicketAgentState) -> dict:
"""等待人工审核"""
timestamp = datetime.now().isoformat()
print("\n" + "="*60)
print("🔴 需要人工审核")
print("="*60)
print(f"会话 ID: {state.get('session_id', 'N/A')}")
print(f"用户 ID: {state.get('user_id', 'N/A')}")
print(f"用户问题: {state.get('query', 'N/A')}")
print(f"问题分类: {state.get('category', 'N/A')}")
print(f"建议方案: {state.get('solution', 'N/A')}")
print(f"风险等级: {state.get('risk_level', 'N/A')}")
print("="*60)
# 等待人工输入
response = input("\n是否批准执行?(yes/no): ").strip().lower()
is_approved = response in ['yes', 'y', '是', '批准']
approval_status = "批准" if is_approved else "拒绝"
audit_entry = f"[{state.get('session_id', 'unknown')}] 人工审核:{approval_status}"
return {
"action_approved": is_approved,
"current_step": "reviewed",
"audit_log": [audit_entry],
"timestamp": [f"审核完成:{timestamp}"]
}
# 6. 执行 API 节点
async def execute_api_node(state: TicketAgentState) -> dict:
"""执行解决方案"""
timestamp = datetime.now().isoformat()
print(f"\n🚀 执行操作:{state.get('solution', 'N/A')}")
# 模拟 API 调用
await asyncio.sleep(2)
result = f"✅ 操作执行成功"
audit_entry = f"[{state.get('session_id', 'unknown')}] 执行完成:{result}"
return {
"progress": "completed",
"current_step": "executed",
"audit_log": [audit_entry],
"timestamp": [f"执行完成:{timestamp}"]
}
# 7. 审计日志节点
async def audit_log_node(state: TicketAgentState) -> dict:
"""生成审计日志"""
timestamp = datetime.now().isoformat()
log_content = "\n".join([
f"=== 工单审计日志 ===",
f"会话 ID: {state.get('session_id', 'N/A')}",
f"用户 ID: {state.get('user_id', 'N/A')}",
f"问题分类: {state.get('category', 'N/A')}",
f"风险等级: {state.get('risk_level', 'N/A')}",
f"执行状态: {state.get('progress', 'N/A')}",
f"",
f"操作记录:",
] + state.get('audit_log', []) + [
f"",
f"时间戳:",
] + state.get('timestamp', []))
print("\n" + log_content)
# 保存到数据库或日志系统
# save_to_database(log_content)
return {
"current_step": "audited",
"timestamp": [f"审计完成:{timestamp}"]
}
# 8. 结束节点
async def end_node(state: TicketAgentState) -> dict:
"""结束会话"""
timestamp = datetime.now().isoformat()
print(f"\n✅ 会话结束 - Session: {state.get('session_id', 'N/A')}")
return {
"progress": "ended",
"current_step": "completed",
"timestamp": [f"会话结束:{timestamp}"]
}
4.2.3 路由函数
def route_based_on_risk_level(state: TicketAgentState) -> str:
"""根据风险等级路由"""
risk_level = state.get('risk_level', 1)
if risk_level >= 3:
print("\n🔴 高风险检测,转入人工审核")
return "human_review"
else:
print(f"\n🟢 低风险 (等级{risk_level}),直接执行")
return "execute_api"
def route_based_on_review_approval(state: TicketAgentState) -> str:
"""根据审核结果路由"""
is_approved = state.get('action_approved', False)
if is_approved:
print("\n✅ 审核通过,执行操作")
return "execute_api"
else:
print("\n❌ 审核拒绝,结束会话")
return "end_session"
4.2.4 图构建与编译
from langgraph.graph import StateGraph
def build_ticket_agent_graph():
"""构建智能客服工单系统图"""
# 1. 创建图实例
graph = StateGraph(TicketAgentState)
# 2. 添加所有节点
graph.add_node("start", start_node)
graph.add_node("classify", classify_query_node)
graph.add_node("generate_solution", generate_solution_node)
graph.add_node("evaluate_risk", evaluate_risk_node)
graph.add_node("human_review", human_review_node)
graph.add_node("execute_api", execute_api_node)
graph.add_node("audit_log", audit_log_node)
graph.add_node("end", end_node)
graph.add_node("end_session", end_node) # 复用 end_node
# 3. 设置入口点
graph.set_entry_point("start")
# 4. 添加顺序边
graph.add_edge("start", "classify")
graph.add_edge("classify", "generate_solution")
graph.add_edge("generate_solution", "evaluate_risk")
# 5. 添加条件边
graph.add_conditional_edges(
"evaluate_risk",
route_based_on_risk_level,
{
"human_review": "human_review",
"execute_api": "execute_api"
}
)
graph.add_conditional_edges(
"human_review",
route_based_on_review_approval,
{
"execute_api": "execute_api",
"end_session": "end_session"
}
)
# 6. 添加后续边
graph.add_edge("execute_api", "audit_log")
graph.add_edge("audit_log", "end")
graph.add_edge("end_session", "end")
# 7. 编译图
app = graph.compile()
return app
# 创建应用实例
app = build_ticket_agent_graph()
4.2.5 执行与测试
import asyncio
async def main():
"""测试智能客服工单系统"""
# 测试用例 1: 低风险问题
print("="*60)
print("测试用例 1: 账户登录问题 (低风险)")
print("="*60)
initial_state_1 = {
"user_id": "user_001",
"query": "我无法登录我的账户,请帮助",
"messages": []
}
final_state_1 = await app.ainvoke(initial_state_1)
print(f"\n最终状态:")
print(f" 会话 ID: {final_state_1['session_id']}")
print(f" 分类: {final_state_1['category']}")
print(f" 风险等级: {final_state_1['risk_level']}")
print(f" 进度: {final_state_1['progress']}")
# 测试用例 2: 高风险问题(需要人工审核)
print("\n" + "="*60)
print("测试用例 2: 投诉问题 (高风险)")
print("="*60)
initial_state_2 = {
"user_id": "user_002",
"query": "我要投诉!这是法律纠纷,我要找媒体曝光!",
"messages": []
}
# 注意:这个测试会触发人工审核,需要输入
# final_state_2 = await app.ainvoke(initial_state_2)
# 测试用例 3: 流式执行
print("\n" + "="*60)
print("测试用例 3: 流式执行演示")
print("="*60)
initial_state_3 = {
"user_id": "user_003",
"query": "我的支付失败了,怎么办?",
"messages": []
}
async for chunk in app.astream(initial_state_3):
for node_name, node_output in chunk.items():
print(f"\n📦 节点 {node_name} 输出:")
for key, value in node_output.items():
if key not in ['audit_log', 'timestamp']: # 简化输出
print(f" {key}: {value}")
if __name__ == "__main__":
asyncio.run(main())
4.3 完整工作流可视化

五、高级特性与最佳实践
5.1 错误处理与重试
from langgraph.errors import GraphInterrupt
import asyncio
async def resilient_node(state: TicketAgentState) -> dict:
"""带错误处理的节点"""
max_retries = 3
retry_count = 0
while retry_count < max_retries:
try:
# 执行可能失败的操作
result = await risky_operation()
return {"result": result, "error": ""}
except Exception as e:
retry_count += 1
print(f"第 {retry_count} 次重试: {str(e)}")
if retry_count >= max_retries:
return {
"result": None,
"error": f"操作失败,已重试{max_retries}次",
"progress": "failed"
}
await asyncio.sleep(1) # 等待 1 秒后重试
5.2 并行执行
from langgraph.graph import END
async def parallel_node(state: TicketAgentState) -> dict:
"""并行执行多个任务"""
# 并行执行多个操作
results = await asyncio.gather(
task1(state),
task2(state),
task3(state),
return_exceptions=True
)
return {
"parallel_results": results,
"current_step": "parallel_completed"
}
# 在图中添加并行分支
graph.add_edge("split_point", "parallel_node")
graph.add_edge("parallel_node", "merge_point")
5.3 状态持久化
from langgraph.checkpoint.memory import MemorySaver
# 启用状态持久化
memory = MemorySaver()
app = graph.compile(checkpointer=memory)
# 保存状态
config = {"configurable": {"thread_id": "thread_123"}}
app.invoke(initial_state, config)
# 恢复状态
checkpoint = memory.get(config)
restored_state = memory.get_state(config)
5.4 性能优化建议
- 批量处理:将多个小节点合并为大节点,减少状态传递开销
- 缓存策略:对重复计算结果进行缓存
- 异步执行:尽可能使用异步节点
- 超时控制:为长时间运行的节点设置超时
- 资源限制:限制并发节点数量
import asyncio
from functools import wraps
def timeout(seconds):
"""节点超时装饰器"""
def decorator(func):
@wraps(func)
async def wrapper(*args, **kwargs):
try:
return await asyncio.wait_for(func(*args, **kwargs), timeout=seconds)
except asyncio.TimeoutError:
return {"error": f"操作超时 ({seconds}秒)", "progress": "timeout"}
return wrapper
return decorator
@timeout(30)
async def slow_node(state: TicketAgentState) -> dict:
"""可能超时的节点"""
await asyncio.sleep(40) # 模拟长时间操作
return {"result": "success"}
六、总结与展望
6.1 核心要点回顾
LangGraph 通过图式编排的理念,为 AI Agent 开发带来了革命性的变化:
- ✅ State:统一的状态管理,避免全局变量混乱
- ✅ Node:模块化的执行单元,便于复用和测试
- ✅ Edge:灵活的流程控制,支持复杂业务逻辑
- ✅ Graph:强大的编排引擎,可视化工作流
6.2 适用场景
LangGraph 特别适合以下场景:
- 多步骤工作流:需要多个步骤串联或并联的业务流程
- 条件分支逻辑:根据中间结果动态调整流程
- 人工介入需求:需要在关键节点进行人工审核
- 状态追踪要求:需要完整的执行轨迹和审计日志
- 复杂决策系统:需要多层判断和路由的智能系统
6.3 未来发展方向
随着 LangGraph 的持续演进,预计将在以下方向取得突破:
- 🚀 更强的可视化能力:实时的执行监控和调试界面
- 🚀 更丰富的节点类型:内置更多预构建的节点模板
- 🚀 更好的性能优化:分布式执行和自动负载均衡
- 🚀 更完善的生态系统:与更多工具和平台的集成
💬 欢迎在评论区交流讨论,如有问题请留言!
更多推荐


所有评论(0)