LangGraph 从入门到精通:AI Agent 编排框架深度解析与实战案例

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

在这里插入图片描述

目录

  1. 为什么需要 LangGraph?
  2. LangGraph 核心概念解析
  3. 与传统框架的对比优势
  4. 实战案例:智能客服工单系统
  5. 高级特性与最佳实践
  6. 总结与展望

一、为什么需要 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 系统架构设计

本案例实现一个完整的智能客服工单系统,具备以下功能:

  1. 自动分类:识别用户问题类型
  2. 方案生成:根据分类生成解决方案
  3. 风险评估:评估操作风险等级
  4. 人工审核:高风险操作需要人工确认
  5. 执行操作:调用 API 执行解决方案
  6. 审计日志:记录完整操作轨迹

在这里插入图片描述

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 性能优化建议

  1. 批量处理:将多个小节点合并为大节点,减少状态传递开销
  2. 缓存策略:对重复计算结果进行缓存
  3. 异步执行:尽可能使用异步节点
  4. 超时控制:为长时间运行的节点设置超时
  5. 资源限制:限制并发节点数量
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 特别适合以下场景:

  1. 多步骤工作流:需要多个步骤串联或并联的业务流程
  2. 条件分支逻辑:根据中间结果动态调整流程
  3. 人工介入需求:需要在关键节点进行人工审核
  4. 状态追踪要求:需要完整的执行轨迹和审计日志
  5. 复杂决策系统:需要多层判断和路由的智能系统

6.3 未来发展方向

随着 LangGraph 的持续演进,预计将在以下方向取得突破:

  • 🚀 更强的可视化能力:实时的执行监控和调试界面
  • 🚀 更丰富的节点类型:内置更多预构建的节点模板
  • 🚀 更好的性能优化:分布式执行和自动负载均衡
  • 🚀 更完善的生态系统:与更多工具和平台的集成

💬 欢迎在评论区交流讨论,如有问题请留言!

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐