LangGraph 0.3 v:外部调用、多轮对话、流式输出、长短期记忆、人机交互、规划逻辑、多 Agent Debug 和监控

 


掌握智能体框架,就是能开发 8 个场景:外部工具调用、多轮对话、流式输出、长短期记忆、人机交互、规划逻辑、多 Agent Debug 和监控

新版本 - 0.3 系列,同时场景更具体,可以快速上手开发
在这里插入图片描述

背景知识

在构建 AI 应用时,我们常常遇到这样的场景:

  • 需要多次调用 AI:一个问题可能需要调用多次大模型才能解决
  • 需要条件判断:根据 AI 的回答决定下一步做什么
  • 需要工具调用:AI 需要查询数据库、搜索网络等
  • 需要循环逻辑:某些任务需要反复执行直到满足条件

传统的 LangChain 是线性的链式调用,难以实现复杂的循环和条件逻辑。

LangGraph 将 AI 应用构建为一个图(Graph)

  • 节点(Node):执行具体任务的函数
  • 边(Edge):连接节点,控制执行流程
  • 状态(State):在节点间传递和共享数据

这种方式让我们可以轻松实现:

用户输入 → AI 分析 → 判断是否需要工具 → 调用工具 → AI 总结 → 返回结果
                ↓                ↓
            直接回答         查询失败重试
  1. 图是整个应用的容器,包含所有节点和边。
from langgraph.graph import StateGraph

builder = StateGraph(State)  # State 定义数据结构
graph = builder.compile()    # 编译后才能运行
  1. 节点是执行具体逻辑的函数,接收状态作为输入,返回更新后的状态。
def my_node(state):
    # 处理逻辑
    result = state["value"] + 1   # 节点函数接收状态字典,对数值进行加1操作
    return {"value": result}      # 返回包含更新后数值的新字典

builder.add_node("my_node", my_node)

节点设计原则:

  1. 单一职责:每个节点只做一件事

  2. 返回部分状态:只返回需要更新的字段

  3. 错误处理:捕获异常并返回错误信息

  4. 边连接节点,控制执行流程。

普通边:固定流向

builder.add_edge(START, "node_a")
builder.add_edge("node_a", "node_b")
builder.add_edge("node_b", END)

条件边:根据条件动态选择下一个节点

def router(state):
    if state["value"] > 10:
        return "node_b"
    else:
        return "node_c"

builder.add_conditional_edges("node_a", router)
  1. 状态是节点间共享的数据结构,使用 TypedDict 定义。
from typing_extensions import TypedDict   
from typing import Annotated
import operator

class State(TypedDict):    # 始终使用 TypedDict 定义 State,对比 dict 优势:类型检查、IDE 提示、减少错误
    messages: Annotated[list, operator.add]  # 自定义的累加消息,operator.add 是自定义的 Reducer函数函数
    count: int                               # 默认模式, 是覆盖旧值

Reducer 就是"更新规则制定者"。

想象你在玩一个游戏:

  • 普通模式:每次更新都覆盖原有分数(默认的 Reducer:覆盖)
  • 累积模式:每次更新都累加到原有分数(自定义 Reducer:累加)

Reducer 就是决定使用哪种模式的"规则"!

Annotated[数据类型, Reducer函数]
          ↓         ↓
      List[str]  operator.add

重点Annotated[list, operator.add] 表示每次更新会追加而不是覆盖

状态设计原则:

  1. 使用 TypedDict:明确字段类型,避免运行时错误
  2. 选择合适的 Reducer
    • operator.add:累加列表
    • 默认:覆盖更新
    • add_messages:智能管理消息
  3. 保持状态简洁:只保存必要的数据

核心概念

  • State 是一个共享的数据结构(通常是字典)
  • 每个节点都可以读取写入 State
  • 节点修改 State 后,会广播给所有其他节点
  • 通过 State 的流动,实现了数据的传递
初始状态: {"x": "10"}[Node 1] 处理并更新 → {"x": "10"}, {"x": "11"}[Node 2] 处理并更新 → {"x": "10"}, {"x": "11"}, {"y": "9"}
    ↓
最终状态: {"x": "10"}, {"x": "11"}, {"y": "9"}

重点理解

  • 节点函数接收完整的 state
  • 只需要返回需要更新的部分
  • LangGraph 会自动合并,不会丢失其他值

 


结构化输出

为什么需要结构化输出?

普通LLM输出:自由文本,格式不固定

结构化输出:固定格式的数据,便于程序处理

用户: "我叫张三,年龄25岁"

❌ 非结构化: "好的,我记住了你的信息"
✅ 结构化: {"name": "张三", "age": 25, "action": "save_to_db"}

使用 Pydantic 定义输出格式

from pydantic import BaseModel

class Output(BaseModel):
    action: str
    confidence: float

structured_llm = llm.with_structured_output(Output)

优势

  • 类型安全
  • 自动验证
  • 易于条件判断

LangGraph支持三种方式:

方式 优点 缺点
Pydantic 类型验证、IDE提示 需要定义类
TypedDict 轻量、类型检查 无运行时验证
JSON Schema 灵活、无需导入 代码冗长
方式1:Pydantic(推荐)

想象你在网上开店:

顾客下单时,你需要检查:

  • 姓名填了吗?是不是文字?
  • 年龄是数字吗?是不是负数?
  • 手机号格式对吗?
  • 收货地址完整吗?

传统做法:你得写一堆 if 判断

def check_order(name, age, phone):
    if not name:
        return "请填写姓名"
    if not isinstance(age, int):
        return "年龄必须是数字"
    if age < 0:
        return "年龄不能是负数"
    if len(phone) != 11:
        return "手机号格式错误"
    # ... 还有100行检查代码

用 Pydantic:告诉它规则,它自动帮你检查

class Order(BaseModel):
    name: str        # 必须是文字
    age: int         # 必须是数字
    phone: str       # 必须是文字

发生了什么?

  1. 你给了两个数据:name="张三"age=25
  2. Pydantic 检查:
    • “张三” 是文字吗?✅ 通过
    • 25 是整数吗?✅ 通过
  3. 创建成功,可以用了

这就是 Pydantic 的作用:你定义规则,它帮你验证数据。

Pydantic 做三件事:

  1. 验证 - 数据对不对(类型、格式、范围)
  2. 转换 - 自动把字符串 “25” 转成数字 25
  3. 整理 - 把数据变成字典或 JSON,方便使用
from pydantic import BaseModel, Field
from typing import Optional

class UserInfo(BaseModel):
    """用户信息"""
    name: str = Field(description="用户姓名")
    age: Optional[int] = Field(description="用户年龄")
    email: str = Field(description="邮箱地址")

# 绑定到LLM,自动验证
structured_llm = llm.with_structured_output(UserInfo)

# 调用
result = structured_llm.invoke("我叫张三,25岁,邮箱zhang@qq.com")
print(result)  # UserInfo(name='张三', age=25, email='zhang@qq.com')

# 在Router中使用
if isinstance(result, UserInfo):
    return "save_to_db"
else:
    return "chat"
方式2:TypedDict
from typing_extensions import Annotated, TypedDict
from typing import Optional

class UserInfo(TypedDict):
    """用户信息"""
    name: Annotated[str, ..., "用户姓名"]
    age: Annotated[Optional[int], None, "用户年龄"]
    email: Annotated[str, ..., "邮箱地址"]

structured_llm = llm.with_structured_output(UserInfo)

result = structured_llm.invoke("我叫张三,25岁,邮箱zhang@qq.com")
print(result)  # {'name': '张三', 'age': 25, 'email': 'zhang@qq.com'}
方式3:JSON Schema
json_schema = {
    "title": "user_info",
    "description": "用户信息",
    "type": "object",
    "properties": {
        "name": {"type": "string", "description": "用户姓名"},
        "age": {"type": "integer", "description": "用户年龄"},
        "email": {"type": "string", "description": "邮箱地址"}
    },
    "required": ["name", "email"]
}

structured_llm = llm.with_structured_output(json_schema)

result = structured_llm.invoke("我叫张三,25岁,邮箱zhang@qq.com")
print(result)  # {'name': '张三', 'age': 25, 'email': 'zhang@qq.com'}
方式 输出类型 优点
Pydantic 对象 类型安全,验证强
TypedDict 字典 轻量,Python原生
JSON Schema 字典 通用,跨语言

推荐使用Pydantic,因为它提供了最好的开发体验和代码可维护性。

 


外部调用

大语言模型(LLM)虽然强大,但存在固有的局限:

局限 说明 示例
知识截止日期 训练数据有时间限制 无法回答今天的天气
无法执行计算 不擅长精确数学运算 计算 1234 × 5678 可能出错
无法访问实时数据 无法联网获取最新信息 无法查询股票实时价格
无法操作系统 无法执行文件操作等 无法创建文件或发送邮件

工具调用(Tool Calling)的价值

通过让 LLM 调用外部工具,我们可以:

  • ✅ 获取实时数据(天气、股票、新闻等)
  • ✅ 执行精确计算(数学运算、数据分析)
  • ✅ 操作外部系统(数据库、API、文件系统)
  • ✅ 扩展 LLM 的能力边界

类比理解

LLM 就像一个聪明的大脑,工具就像手和眼睛。大脑能思考但不能直接感知世界和操作物体,需要通过感官器官(工具)来完成实际任务。

核心流程:

用户提问
    ↓
LLM 判断是否需要工具
    ↓
【需要工具】→ LLM 生成工具调用请求
    ↓
执行工具(查询天气/计算/搜索等)
    ↓
工具返回结果
    ↓
LLM 整合结果生成最终回答
    ↓
返回给用户
🔹 工具定义(Tool Definition)

告诉 LLM 有哪些工具可用

from langchain_core.tools import tool

@tool
def get_weather(city: str) -> str:
    """查询指定城市的天气情况

    Args:
        city: 城市名称,如"北京"、"上海"

    Returns:
        天气描述字符串
    """
    # 模拟天气查询
    weather_data = {
        "北京": "晴天,温度 25°C",
        "上海": "多云,温度 28°C",
        "广州": "阵雨,温度 30°C"
    }
    return weather_data.get(city, "未知城市")

关键点

  • 使用 @tool 装饰器
  • docstring (三引号内的文档字符串)非常重要:LLM 通过它理解工具的用途
  • 参数需要类型注解
@tool 装饰器
@tool(
    name: str = None,           # 工具名称,默认使用函数名
    description: str = None,    # 工具描述,默认使用 docstring
    return_direct: bool = False # 是否直接返回工具结果(不再给 LLM)
)
🔹 工具绑定(Tool Binding)

将工具绑定到 LLM

from langchain_openai import ChatOpenAI

# 初始化大语言模型实例,指定使用 GPT-4 模型
llm = ChatOpenAI(model="gpt-4")

# 将天气查询工具绑定到 LLM,使模型能够识别并调用该工具
llm_with_tools = llm.bind_tools([get_weather])
🔹 工具节点(ToolNode)

在 LangGraph 中执行工具

from langgraph.prebuilt import ToolNode

# 创建工具执行节点,专门用于在 LangGraph 工作流中运行工具函数
# 当工作流需要执行工具调用时,请求会被路由到此节点
tool_node = ToolNode([get_weather])

核心概念解释:

工具绑定(Tool Binding)

  • 目的:让 LLM 知道有哪些工具可用
  • 作用:模型学会在适当场景决定是否调用工具
  • 结果:模型输出中包含工具调用请求

工具节点(Tool Node)

  • 目的:在工作流中实际执行工具调用
  • 作用:接收工具调用请求,运行对应函数并返回结果
  • 结果:获得工具执行的实际数据(如天气信息)

工作流程

  1. 绑定阶段:LLM 学习工具知识 → 生成调用意图
  2. 执行阶段:ToolNode 接收意图 → 执行工具函数 → 返回真实数据
llm.bind_tools()

作用:告诉 LLM 有哪些工具可以使用

核心价值

  • LLM 会在需要时主动调用工具
  • 支持多个工具同时绑定
  • 可控制工具调用策略
llm.bind_tools(
    tools: list[BaseTool],      # 工具列表
    tool_choice: str = "auto"   # 工具选择策略
)

tool_choice 参数

  • "auto"(默认):LLM 自主决定是否调用工具
  • "required":强制 LLM 必须调用工具
  • "none":禁止 LLM 调用工具
  • 具体工具名:强制调用指定工具
ToolNode

作用:在 LangGraph 图中执行工具调用

核心价值

  • 自动解析 LLM 的工具调用请求
  • 执行对应的工具函数
  • 将结果包装成 ToolMessage 返回给 LLM
ToolNode 返回值
{
    "messages": [ToolMessage(
        content="工具执行结果",
        tool_call_id="call_xxx",
        name="get_weather"
    )]
}

ToolMessage 字段

  • content:工具返回的内容
  • tool_call_id:对应 LLM 工具调用请求的 ID
  • name:被调用的工具名称
用户输入: "北京今天天气怎么样?"[agent 节点]
- LLM 分析问题:需要查询北京天气
- 生成工具调用请求:
  {
    "tool": "get_weather",
    "arguments": {"city": "北京"}
  }[tools_condition 判断] → 检测到工具调用 → 路由到 "tools"[tools 节点]
- 执行 get_weather("北京")
- 返回: "晴天,温度 25°C,空气质量优"[返回到 agent 节点]
- LLM 接收工具结果
- 生成自然语言回答: "北京今天是晴天..."[tools_condition 判断] → 没有工具调用 → 路由到 END
    ↓
返回最终结果
完整案例:多工具协同

场景:旅行助手

我们扩展助手功能,支持:

  • 查询天气
  • 计算器(计算旅行预算)
  • 查询交通信息
定义多个工具
@tool
def get_weather(city: str) -> str:
    """查询城市天气"""
    # (同上)
    pass

@tool
def calculator(expression: str) -> str:
    """执行数学计算

    Args:
        expression: 数学表达式,如 "100 * 3 + 50"

    Returns:
        计算结果
    """
    try:
        result = eval(expression)
        return f"计算结果: {result}"
    except Exception as e:
        return f"计算错误: {str(e)}"

@tool
def get_transportation(start_city: str, end_city: str) -> str:
    """查询两个城市之间的交通方式

    Args:
        start_city: 出发城市
        end_city: 目的地城市

    Returns:
        交通方式建议
    """
    return f"从{start_city}{end_city}可以选择:\n" \
           f"- 高铁:约 3 小时\n" \
           f"- 飞机:约 1.5 小时\n" \
           f"- 大巴:约 5 小时"
构建多工具 Agent
from langgraph.graph import StateGraph, START, END
from langgraph.prebuilt import ToolNode, tools_condition

# 所有工具
tools = [get_weather, calculator, get_transportation]

# 绑定到 LLM
llm_with_tools = llm.bind_tools(tools)

# 定义 agent 节点
def agent(state: AgentState):
    response = llm_with_tools.invoke(state["messages"])
    return {"messages": [response]}

# 构建图
builder = StateGraph(AgentState)
builder.add_node("agent", agent)
builder.add_node("tools", ToolNode(tools))

builder.add_edge(START, "agent")
builder.add_conditional_edges("agent", tools_condition, {"tools": "tools", END: END})
builder.add_edge("tools", "agent")

graph = builder.compile()
测试多工具调用
response = graph.invoke({
    "messages": [HumanMessage(content="""
        我计划从北京去上海旅行,请帮我:
        1. 查询上海的天气
        2. 推荐交通方式
        3. 如果高铁票价 500 元,3 个人一共多少钱?
    """)]
})

print(response["messages"][-1].content)

LLM 会自动调用多个工具

  1. get_weather("上海") → 获取天气
  2. get_transportation("北京", "上海") → 获取交通信息
  3. calculator("500 * 3") → 计算总价

最终返回整合的回答。

工具调用技巧:
  1. 工具描述要清晰:AI 根据描述选择工具
@tool
def search_web(query: str) -> str:
    """在互联网上搜索信息。适用于需要最新信息的查询。"""
    pass
  1. 参数类型要明确:使用 Pydantic 定义参数
class SearchInput(BaseModel):
    query: str = Field(description="搜索关键词")

@tool(args_schema=SearchInput)
def search_web(query: str) -> str:
    pass
  1. 处理工具失败:返回有意义的错误信息
@tool
def search_web(query: str) -> str:
    try:
        return search_api(query)
    except Exception as e:
        return f"搜索失败: {e}"

Q1: LLM 什么时候会调用工具?

A: LLM 会在认为需要外部信息时调用工具。关键因素:

  • 工具的 docstring 描述是否清晰
  • 用户问题是否明确需要该工具
  • 工具名称是否语义明确

Q2: 如何让 LLM 更准确地调用工具?

A:

  1. 写清晰的 docstring,说明工具的用途和参数
  2. 给工具起语义化的名字
  3. 在系统提示中引导 LLM 使用工具
system_message = """
你是一个智能助手。当用户询问天气时,请使用 get_weather 工具查询。
"""

Q3: 工具调用失败怎么办?

A: 在工具函数中添加错误处理:

@tool
def get_weather(city: str) -> str:
    """查询天气"""
    try:
        # 实际的 API 调用
        result = weather_api.query(city)
        return result
    except Exception as e:
        # 返回友好的错误信息
        return f"抱歉,查询{city}的天气时出错:{str(e)}"

Q4: 如何限制工具调用次数?

A: 使用递归限制参数:

result = graph.invoke(
    {"messages": [HumanMessage(content="...")]},
    {"recursion_limit": 10}  # 最多执行 10 次节点
)
  1. 详细的工具文档
@tool
def search_database(query: str, limit: int = 5) -> str:
    """在数据库中搜索相关记录

    使用场景:
    - 用户询问历史数据时使用
    - 需要查找特定信息时使用

    Args:
        query: 搜索关键词,支持模糊匹配
        limit: 返回结果数量,默认 5 条,最多 20 条

    Returns:
        JSON 格式的搜索结果

    Examples:
        >>> search_database("用户反馈", limit=3)
        >>> search_database("订单查询")
    """
    # 实现...
  1. 参数验证
@tool
def get_weather(city: str) -> str:
    """查询天气"""
    if not city or len(city) == 0:
        return "错误:城市名称不能为空"

    if len(city) > 20:
        return "错误:城市名称过长"

    # 正常逻辑...
  1. 返回结构化数据
import json

@tool
def get_stock_price(symbol: str) -> str:
    """查询股票价格"""
    data = {
        "symbol": symbol,
        "price": 150.25,
        "change": +2.5,
        "timestamp": "2024-01-09 10:30:00"
    }
    # 返回 JSON 字符串,LLM 更容易解析
    return json.dumps(data, ensure_ascii=False)
❌ DON’T(避免)
  1. ❌ 模糊的工具描述
@tool
def do_something(x):
    """做一些事情"""  # 太模糊!
    pass
  1. ❌ 工具内部直接返回异常
@tool
def get_weather(city: str) -> str:
    result = api.call(city)  # 如果 API 挂了,会抛异常
    return result            # LLM 会看到 Python 堆栈跟踪
  1. ❌ 工具执行耗时操作
@tool
def process_large_file(file_path: str) -> str:
    """处理大文件"""
    # 可能需要几分钟,会导致超时
    time.sleep(300)
    return "完成"

 


多轮对话

价值

  • ✅ 连贯的对话体验
  • ✅ 上下文相关的回答
  • ✅ 支持复杂的多步骤任务

在构建聊天机器人时,我们需要:

  • 📝 保存用户的每一条提问
  • 🤖 记录 AI 的每一次回复
  • 🔄 维护完整的对话历史

这就需要一个专门的消息管理机制!

消息类型

LangChain 定义了多种消息类型:

消息类型 作用 示例
HumanMessage 用户输入 HumanMessage(content="你好")
AIMessage AI 回复 AIMessage(content="你好!有什么可以帮你的?")
SystemMessage 系统提示 SystemMessage(content="你是一个友好的助手")
ToolMessage 工具返回 ToolMessage(content="天气:晴天", tool_call_id="xxx")
from langchain_core.messages import (
    HumanMessage,
    AIMessage,
    SystemMessage,
    ToolMessage
)

# 创建消息
user_msg = HumanMessage(content="你好")
ai_msg = AIMessage(content="你好!有什么可以帮你的?")
system_msg = SystemMessage(content="你是一个友好的助手")
messages: Annotated[Sequence[BaseMessage], operator.add]

解释

  • Sequence[BaseMessage]:消息列表的类型
  • Annotated[..., operator.add]:合并策略
    • operator.add:追加模式(类似列表的 + 操作)
    • 新消息会追加到旧消息后面

构建带对话历史的聊天机器人:

运行效果

用户: 你好
AI: 你好!有什么可以帮助你的吗?

用户: 介绍一下你自己
AI: 我是一个由OpenAI开发的人工智能助手...

用户: 退出
# 导入必要的类型注解模块
from typing import Annotated
# 导入用于定义类型化字典的模块
from typing_extensions import TypedDict
# 导入LangGraph的图构建相关类
from langgraph.graph import StateGraph, START, END
# 导入消息处理工具函数
from langgraph.graph.message import add_messages
# 导入OpenAI聊天模型
from langchain_openai import ChatOpenAI

# Step 1: 定义状态类
class State(TypedDict):
    # 定义消息列表,使用Annotated添加元数据
    messages: Annotated[list, add_messages]
    #                         ^^^^^^^^^^^
    # LangGraph提供的智能消息处理器,自动处理消息的合并和更新

# Step 2: 创建图和模型
# 初始化状态图构建器,传入状态类定义
graph_builder = StateGraph(State)
# 创建OpenAI聊天模型实例,使用gpt-4o模型
llm = ChatOpenAI(model="gpt-4o")

# Step 3: 定义聊天节点函数
def chatbot(state: State):
    """调用大模型,返回 AI 回复"""
    # 调用大模型处理当前对话消息
    return {"messages": [llm.invoke(state["messages"])]}

# Step 4: 构建图结构
# 向图中添加名为"chatbot"的节点,关联到chatbot函数
graph_builder.add_node("chatbot", chatbot)
# 添加从开始节点到chatbot节点的边
graph_builder.add_edge(START, "chatbot")
# 添加从chatbot节点到结束节点的边
graph_builder.add_edge("chatbot", END)

# 编译图,生成可执行的图实例
graph = graph_builder.compile()

交互工作流:

# 定义流式输出图更新结果的函数
def stream_graph_updates(user_input: str):
    # 遍历图流中产生的事件
    for event in graph.stream({"messages": [("user", user_input)]}):
        # 遍历事件中的每个值
        for value in event.values():
            # 打印AI的最新回复内容
            print("AI:", value["messages"][-1].content)

# 使用示例
# 创建无限循环以实现持续对话
while True:
    # 获取用户输入
    user_input = input("用户: ")
    # 检查是否为退出命令
    if user_input.lower() == "退出":
        # 跳出循环,结束程序
        break
    # 调用函数处理用户输入并获取AI回复
    stream_graph_updates(user_input)

深入 add_messages (智能合并消息列表)源码:

from langgraph.graph.message import add_messages

# 情况1:不同 ID → 追加
msgs1 = [HumanMessage(content="你好", id="1")]
msgs2 = [AIMessage(content="你好!", id="2")]
result = add_messages(msgs1, msgs2)
# 结果: [HumanMessage(id="1"), AIMessage(id="2")]

# 情况2:相同 ID → 覆盖
msgs1 = [HumanMessage(content="你好", id="1")]
msgs2 = [HumanMessage(content="你好呀", id="1")]  # 同一个 ID
result = add_messages(msgs1, msgs2)
# 结果: [HumanMessage(content="你好呀", id="1")]  ← 被更新了!

3. 访问消息历史

def analyze_conversation(state: ChatState):
    """分析对话历史"""
    messages = state["messages"]

    # 获取最后一条用户消息
    last_user_message = None
    for msg in reversed(messages):
        if isinstance(msg, HumanMessage):
            last_user_message = msg
            break

    # 统计消息数量
    user_count = sum(1 for m in messages if isinstance(m, HumanMessage))
    ai_count = sum(1 for m in messages if isinstance(m, AIMessage))

    print(f"用户消息: {user_count}, AI 消息: {ai_count}")

对比 operator.add

特性 operator.add add_messages
追加新消息 ✅ 支持 ✅ 支持
更新现有消息 ❌ 不支持 ✅ 支持(通过 ID)
删除消息 ❌ 不支持 ✅ 支持(RemoveMessage)
智能去重 ❌ 不支持 ✅ 支持
适用场景 简单累加 复杂对话管理

推荐使用场景

  • 🔹 简单数据累加 → operator.add
  • 🔸 对话历史管理 → add_messages
  • 🔸 需要消息更新 → add_messages
  • 🔸 人机交互应用 → add_messages

MessageGraph:专为对话设计:

from langgraph.graph.message import MessageGraph

# 创建专门处理消息的图
builder = MessageGraph()

# 添加节点(自动处理消息列表)
builder.add_node("chatbot", lambda state: [("assistant", "你好!")])

# 编译和运行
graph = builder.compile()
result = graph.invoke([("user", "你好")])
# 结果: [HumanMessage(...), AIMessage(...)]

特点

  • ✅ 自动使用 add_messages 管理状态
  • ✅ 状态就是消息列表,简单直观
  • ✅ 适合纯对话应用
  • ❌ 无法存储其他复杂数据

StateGraph:通用灵活:

from langgraph.graph import StateGraph

class State(TypedDict):
    messages: Annotated[list, add_messages]
    user_info: dict  # 可以存储其他数据!
    context: str
    step_count: int

builder = StateGraph(State)

特点

  • ✅ 支持复杂数据结构
  • ✅ 可自定义多个键和 Reducer
  • ✅ 适合复杂业务逻辑
  • ⚠️ 需要手动配置

需求:构建一个聊天助手,需要:

  1. 保存对话历史
  2. 提取用户信息
  3. 将 AI 回复转换为 JSON 格式
from langchain_openai import ChatOpenAI
from langchain_core.messages import SystemMessage, HumanMessage, AIMessage
from langgraph.graph import StateGraph, END
from typing import Annotated, TypedDict, List
from langgraph.graph.message import add_messages
import operator

# Step 1: 定义 State
class State(TypedDict):
    messages: Annotated[List, add_messages]

# Step 2: 初始化模型
llm = ChatOpenAI(model="gpt-4o")

# Step 3: 定义节点1 - 聊天
def chat_with_model(state):
    """与模型对话"""
    messages = state['messages']
    response = llm.invoke(messages)
    return {"messages": [response]}

# Step 4: 定义节点2 - 转换为 JSON
def convert_to_json(state):
    """将 AI 回复转换为 JSON 格式"""
    EXTRACTION_PROMPT = """
    你是数据提取专家,请将以下文本提取为 JSON 格式。
    只输出 JSON,不要其他内容。
    """

    last_message = state['messages'][-1].content
    messages = [
        SystemMessage(content=EXTRACTION_PROMPT),
        HumanMessage(content=last_message)
    ]

    response = llm.invoke(messages)
    return {"messages": [response]}

# Step 5: 构建图
builder = StateGraph(State)
builder.add_node("chat", chat_with_model)
builder.add_node("convert", convert_to_json)

builder.set_entry_point("chat")
builder.add_edge("chat", "convert")
builder.add_edge("convert", END)

graph = builder.compile()

# Step 6: 运行
query = "你好,请介绍一下你自己"
result = graph.invoke({"messages": [HumanMessage(content=query)]})

# 查看结果
print(result["messages"][-1].content)

执行流程

用户输入: "你好,请介绍一下你自己"
    ↓
chat 节点: 生成自我介绍
    ↓
convert 节点: 将介绍转为 JSON
    ↓
输出:
{
  "entity": "OpenAI",
  "role": "AI Assistant",
  "capabilities": [...]
}

Reducer 是状态更新规则,定义如何合并新旧值:

1. 计数器 Reducer

counts: Annotated[int, lambda x, y: x + y]
# 旧值x=2, 新值y=3 → 结果=5
# 用途:统计调用次数、计数等

2. 列表追加 Reducer

items: Annotated[List, operator.add]
# 旧列表x=[1,2], 新列表y=[3,4] → 结果=[1,2,3,4]
# 用途:收集多个节点的输出

3. 字典合并 Reducer

config: Annotated[dict, lambda x, y: {**x, **y}]
# 旧字典x={"a":1}, 新字典y={"b":2} → 结果={"a":1, "b":2}
# 用途:合并配置参数

4. 消息专用 Reducer

messages: Annotated[list, add_messages]
# 智能处理:去重、排序、合并对话消息
# 用途:管理对话历史

Reducer 函数

  • 控制 State 的更新规则
  • 默认:覆盖
  • operator.add:累加
  • add_messages:智能消息管理

问:如何实现多轮对话?

答:使用 add_messages 管理历史消息:

class State(TypedDict):
    messages: Annotated[list, add_messages]

# 每次调用会累加消息
graph.invoke({"messages": [HumanMessage(content="第一轮")]})
graph.invoke({"messages": [HumanMessage(content="第二轮")]})
实践技巧

Q2: 如何在多个用户之间隔离对话?

A: 使用 thread_id 区分不同用户的会话(第四章详解):

# 用户 A 的对话
config_a = {"configurable": {"thread_id": "user_a"}}
graph.invoke({"messages": [...]}, config_a)

# 用户 B 的对话
config_b = {"configurable": {"thread_id": "user_b"}}
graph.invoke({"messages": [...]}, config_b)

Q3: MessagesState 可以和其他字段一起使用吗?

A: 可以!只需继承并添加字段:

class MyState(MessagesState):
    messages: ...  # 继承自 MessagesState
    user_name: str  # 自定义字段
    session_id: str  # 自定义字段
    context: dict  # 自定义字段

Q4: 如何修改历史消息?

A: 直接修改会破坏追加逻辑,建议返回新的完整列表:

def modify_last_message(state: MessagesState):
    """修改最后一条消息"""
    messages = state["messages"].copy()

    # 修改最后一条
    messages[-1].content = "修改后的内容"

    # 返回完整列表(不使用追加)
    return {"messages": messages}

 


流式输出

在没有流式输出时,用户必须等待整个处理完成:

问题

  • ❌ 用户不知道系统在做什么
  • ❌ 无法感知进度
  • ❌ 等待时间长,体验差
  • ❌ 无法提前看到部分结果

价值

  • ✅ 实时反馈,用户知道系统在工作
  • ✅ 提前看到部分结果
  • ✅ 更好的用户体验
  • ✅ 方便调试和监控

类比理解

就像观看视频加载进度条,虽然总时间一样,但有进度反馈会让等待更轻松。

三种 Stream Mode

LangGraph 提供了三种流式输出模式,适用于不同场景。

参数是什么意思?

stream() 方法签名
graph.stream(
    input: dict,                    # 初始输入状态
    config: dict = None,            # 可选配置(如 thread_id)
    stream_mode: str = "values",    # 流式模式
    **kwargs
) -> Iterator[dict]

stream_mode 参数详解

模式 输出内容 适用场景
"values" 完整状态字典 需要看到每个节点后的完整状态
"updates" 节点的更新内容 只关心新增的部分,节省带宽
"messages" 新增的消息对象 聊天机器人,只关心对话内容

config 参数示例

config = {
    "configurable": {
        "thread_id": "user_123",  # 会话标识
    },
    "recursion_limit": 20  # 最大递归深度
}

for chunk in graph.stream(input_data, config=config):
    print(chunk)
2.1 stream_mode=“values”(完整状态流)

特点:输出每个节点执行后的完整状态

for chunk in graph.stream(input, stream_mode="values"):
    print(chunk)  # chunk 是完整的状态字典

输出格式

# 第1个 chunk(节点1执行后)
{"messages": [HumanMessage(...), AIMessage(...)], "user_name": "张三"}

# 第2个 chunk(节点2执行后)
{"messages": [HumanMessage(...), AIMessage(...), AIMessage(...)], "user_name": "张三"}

# 第3个 chunk(节点3执行后)
{"messages": [...], "user_name": "张三", "result": "完成"}

适用场景

  • ✅ 需要看到每个节点后的完整状态
  • ✅ 调试时追踪状态变化
  • ✅ 监控整个执行过程
2.2 stream_mode=“updates”(增量更新流)

特点:只输出每个节点的更新内容

for chunk in graph.stream(input, stream_mode="updates"):
    print(chunk)  # chunk 只包含新增/修改的部分

输出格式

# 第1个 chunk(节点1的更新)
{"node_1": {"messages": [AIMessage(...)]}}

# 第2个 chunk(节点2的更新)
{"node_2": {"messages": [AIMessage(...)], "result": "部分完成"}}

# 第3个 chunk(节点3的更新)
{"node_3": {"result": "完成"}}

适用场景

  • ✅ 只关心新增的内容
  • ✅ 需要知道是哪个节点的输出
  • ✅ 减少数据传输量
2.3 stream_mode=“messages”(消息流)

特点:只输出新的消息(适用于对话场景)

for chunk in graph.stream(input, stream_mode="messages"):
    print(chunk)  # chunk 是单个新消息

输出格式

# 第1个 chunk
AIMessage(content="你好!")

# 第2个 chunk
AIMessage(content="我正在查询天气...")

# 第3个 chunk
ToolMessage(content="北京:晴天,25°C")

# 第4个 chunk
AIMessage(content="北京今天是晴天,温度25°C。")

适用场景

  • ✅ 聊天机器人
  • ✅ 只关心对话内容
  • ✅ 简化消息处理

Q1: stream() 和 invoke() 的区别是什么?

A:

方法 特点 适用场景
invoke() 等待全部完成,一次性返回 不需要实时反馈的场景
stream() 实时返回中间结果 需要展示进度,提升体验
# invoke:等待完成
result = graph.invoke(input)  # 阻塞 10 秒
print(result)  # 一次性输出

# stream:实时输出
for chunk in graph.stream(input):  # 每个节点都输出
    print(chunk)

Q2: 如何只显示 AI 的最终回复,不显示中间工具调用?

A: 使用 stream_mode="messages" 并过滤:

for chunk in graph.stream(input, stream_mode="messages"):
    # 只显示 AI 的文本回复
    if isinstance(chunk, AIMessage) and chunk.content:
        print(chunk.content)

Q3: 流式输出会影响性能吗?

A: 不会。流式输出只是改变了数据的返回方式,不影响执行速度。

Q4: 如何在 Web 应用中使用流式输出?

A: 使用 Server-Sent Events (SSE):

from flask import Flask, Response

app = Flask(__name__)

@app.route('/stream')
def stream():
    def generate():
        for chunk in graph.stream(input, stream_mode="messages"):
            yield f"data: {chunk.content}\n\n"

    return Response(generate(), mimetype='text/event-stream')

 


短期记忆

在之前的章节中,我们构建的图都有一个明显的局限:每次执行后,图会"忘记"之前的对话

问题:

  • ❌ 无法维护对话历史
  • ❌ 每次执行都是全新的开始
  • ❌ 无法实现多轮对话
  • ❌ 用户体验差

State 状态的局限性

你可能会想:“不是有 State 状态吗?为什么不能记住?”

关键理解:

State 状态模式只在单次运行期间维护消息状态。每次运行完成后,状态会重置到初始状态。

State 的设计目的:

  • ✅ 在单次运行期间传递节点间的数据
  • ❌ 不是为了跨运行保存历史

类比理解:

State 就像大脑的"工作记忆"(只在做一件事时起作用),而我们需要的是"长期记忆"(可以记住昨天、上周甚至更早的事情)。

LangGraph 将记忆分为两类:

记忆类型 作用范围 实现方式 适用场景
短期记忆 单个会话线程内 Checkpointer 记住当前对话的历史
长期记忆 跨会话线程 Store 记住用户偏好、历史交互
🔹 Checkpointer(检查点)

作用:在每个超级步骤(super-step)后保存图状态的快照

工作原理:

执行节点1 → 保存检查点1 → 执行节点2 → 保存检查点2 → ...

关键点:

  • 每次状态更新都会创建一个检查点
  • 检查点保存在 thread(线程)中
  • 可以通过 thread_id 访问特定会话的历史
🔹 Thread(线程)

作用:唯一标识一个对话会话

类比理解:

Thread 就像微信中的不同聊天窗口,每个窗口(thread)有自己的对话历史,互不干扰。

# 线程1:与用户A的对话
config_a = {"configurable": {"thread_id": "user_a"}}

# 线程2:与用户B的对话
config_b = {"configurable": {"thread_id": "user_b"}}
🔹 Store(仓库)

作用:跨线程存储和检索信息

使用场景:

  • 用户偏好设置
  • 知识库
  • 所有线程的全局配置

❓ 问题1:Checkpointer这个模块是干什么用的?

Checkpointer 是 LangGraph 实现短期记忆的核心机制。

核心功能:

  1. 自动保存状态:每个节点执行后自动创建检查点
  2. 状态恢复:根据 thread_id 恢复之前的对话历史
  3. 会话隔离:不同 thread_id 的对话互不干扰

核心价值:

将无状态的图转变为有记忆的对话系统

❓ 问题2:最常用的方法是什么?

1. MemorySaver:内存存储
from langgraph.checkpoint.memory import MemorySaver

# 创建检查点
memory = MemorySaver()

# 编译图时传入
graph = builder.compile(checkpointer=memory)

# 指定 thread_id 运行
config = {"configurable": {"thread_id": "conversation_1"}}
result = graph.invoke(input_data, config)

特点:

  • 存储在内存中
  • 程序重启后丢失
  • 适合演示和测试
2. SqliteSaver:持久化存储
from langgraph.checkpoint.sqlite import SqliteSaver

# 方式1:内存存储(同 MemorySaver)
with SqliteSaver.from_conn_string(":memory:") as memory:
    graph = builder.compile(checkpointer=memory)

# 方式2:持久化到文件
with SqliteSaver.from_conn_string("checkpoints.db") as memory:
    graph = builder.compile(checkpointer=memory)

特点:

  • 可以持久化到数据库
  • 程序重启后数据仍然存在
  • 适合生产环境
3. 跨单元使用的最佳实践
from contextlib import ExitStack

# 创建持久化的 checkpointer
stack = ExitStack()
checkpointer = stack.enter_context(
    SqliteSaver.from_conn_string(":memory:")
)

# 编译图
graph = builder.compile(checkpointer=checkpointer)

# 可以在多个单元中使用
config = {"configurable": {"thread_id": "1"}}
graph.invoke(input1, config)
graph.invoke(input2, config)  # 记得之前的对话

# 使用完后关闭
stack.close()
4. 异步版本
from contextlib import AsyncExitStack
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver

# 异步环境下使用
stack = AsyncExitStack()
memory = await stack.enter_async_context(
    AsyncSqliteSaver.from_conn_string(":memory:")
)

graph = builder.compile(checkpointer=memory)

# 异步执行
async for chunk in graph.astream(input, config):
    print(chunk)

# 关闭
await stack.aclose()

❓ 问题3:参数是什么意思?

compile(checkpointer=…)
graph.compile(
    checkpointer: BaseCheckpointSaver = None
)

checkpointer 参数:

  • None(默认):不保存检查点,无记忆功能
  • MemorySaver():内存存储
  • SqliteSaver.from_conn_string(...):SQLite 持久化
config 参数
config = {
    "configurable": {
        "thread_id": str,          # 必需,会话唯一标识
        "checkpoint_id": str,      # 可选,特定检查点 ID
        "checkpoint_ns": str       # 可选,子图命名空间
    },
    "recursion_limit": int         # 可选,最大递归深度
}

thread_id:

  • 必需参数(使用 checkpointer 时)
  • 用于隔离不同会话
  • 相同 thread_id 共享历史

示例:

# 用户 A 的对话
config_a = {"configurable": {"thread_id": "user_a"}}
graph.invoke({"messages": [HumanMessage("我叫张三")]}, config_a)

# 用户 B 的对话(独立的历史)
config_b = {"configurable": {"thread_id": "user_b"}}
graph.invoke({"messages": [HumanMessage("我叫李四")]}, config_b)

# 继续用户 A 的对话
graph.invoke({"messages": [HumanMessage("我叫什么?")]}, config_a)
# 输出:你叫张三。

实战案例:带记忆的智能客服

需求分析

构建一个智能客服系统,能够:

  1. 记住用户的姓名和问题
  2. 在同一会话中保持上下文
  3. 支持多个用户同时对话(不同 thread_id)

4.2 完整实现

Step 1:使用 MemorySaver
from langgraph.graph import StateGraph, MessagesState, START, END
from langgraph.checkpoint.memory import MemorySaver
from langchain_core.messages import HumanMessage
from langchain_openai import ChatOpenAI

# 初始化 LLM
llm = ChatOpenAI(model="gpt-4")

# 定义状态(使用 MessagesState)
class CustomerServiceState(MessagesState):
    pass

# 定义客服节点
def customer_service_agent(state: CustomerServiceState):
    """处理用户消息"""
    response = llm.invoke(state["messages"])
    return {"messages": [response]}

# 构建图
builder = StateGraph(CustomerServiceState)
builder.add_node("agent", customer_service_agent)
builder.add_edge(START, "agent")
builder.add_edge("agent", END)

# 添加记忆功能
memory = MemorySaver()
graph = builder.compile(checkpointer=memory)
Step 2:多轮对话测试
# 配置 thread_id
config = {"configurable": {"thread_id": "customer_001"}}

# 第1轮:自我介绍
result = graph.invoke(
    {"messages": [HumanMessage(content="你好,我叫张三")]},
    config
)
print("AI:", result["messages"][-1].content)
# 输出:你好,张三!很高兴为您服务。

# 第2轮:询问问题
result = graph.invoke(
    {"messages": [HumanMessage(content="我想退款")]},
    config
)
print("AI:", result["messages"][-1].content)
# 输出:好的,张三,我帮您处理退款...

# 第3轮:验证记忆
result = graph.invoke(
    {"messages": [HumanMessage(content="我叫什么名字?")]},
    config
)
print("AI:", result["messages"][-1].content)
# 输出:您叫张三。
Step 3:多用户场景
# 用户 A
config_a = {"configurable": {"thread_id": "customer_001"}}
graph.invoke(
    {"messages": [HumanMessage("我是张三,想退款")]},
    config_a
)

# 用户 B(独立的会话)
config_b = {"configurable": {"thread_id": "customer_002"}}
graph.invoke(
    {"messages": [HumanMessage("我是李四,咨询产品")]},
    config_b
)

# 继续用户 A 的对话
graph.invoke(
    {"messages": [HumanMessage("我的退款处理好了吗?")]},
    config_a
)
# AI 知道这是张三,知道他之前要退款

持久化存储:SqliteSaver

为什么需要持久化?

MemorySaver 的局限:

  • 数据只存在于内存中
  • 程序重启后数据丢失
  • 无法跨进程共享

SqliteSaver 的优势:

  • 数据保存在数据库文件中
  • 程序重启后数据仍然存在
  • 可以查询和管理历史数据

5.2 使用 SqliteSaver

方式1:内存模式
from langgraph.checkpoint.sqlite import SqliteSaver

with SqliteSaver.from_conn_string(":memory:") as memory:
    graph = builder.compile(checkpointer=memory)

    config = {"configurable": {"thread_id": "1"}}
    graph.invoke(input, config)
方式2:文件持久化
with SqliteSaver.from_conn_string("chatbot.db") as memory:
    graph = builder.compile(checkpointer=memory)

    config = {"configurable": {"thread_id": "1"}}
    graph.invoke(input, config)

文件持久化的好处:

  • chatbot.db 文件会保存所有检查点
  • 程序重启后仍然可以访问历史
  • 可以用 SQL 工具查看数据
查看数据库内容
import sqlite3

# 连接数据库
conn = sqlite3.connect("chatbot.db")
cursor = conn.cursor()

# 查看所有表
cursor.execute("SELECT name FROM sqlite_master WHERE type='table';")
print(cursor.fetchall())
# 输出:[('checkpoints',), ('writes',)]

# 查询检查点数据
cursor.execute("SELECT * FROM checkpoints;")
for row in cursor.fetchall():
    print(row)

conn.close()

 


长期记忆

短期记忆实现

  • MemorySaver:内存存储
  • SqliteSaver:持久化存储
  • thread_id:会话管理

长期记忆实现

  • InMemoryStore:跨线程存储
  • Namespace:数据组织
  • 结合 Checkpointer 和 Store

Checkpointer 的局限

虽然 Checkpointer 实现了短期记忆,但它有一个限制:无法跨线程共享信息

# 线程1
config1 = {"configurable": {"thread_id": "1"}}
graph.invoke({"messages": [HumanMessage("我叫张三")]}, config1)

# 线程2(无法访问线程1的信息)
config2 = {"configurable": {"thread_id": "2"}}
graph.invoke({"messages": [HumanMessage("我叫什么?")]}, config2)
# AI 不知道,因为不同的 thread_id

需要的场景:

  • 用户偏好设置(所有会话都适用)
  • 知识库(所有用户共享)
  • 全局配置

6.2 InMemoryStore

InMemoryStore 提供跨线程的存储能力。

核心概念

Namespace(命名空间):

  • 类似文件夹,组织数据
  • 格式:(user_id, "memories")
  • 支持层级结构

Key(键):

  • 每条记忆的唯一标识
  • 通常使用 UUID
基本用法
from langgraph.store.memory import InMemoryStore
import uuid

# 创建存储
store = InMemoryStore()

# 定义命名空间
user_id = "user_123"
namespace = (user_id, "memories")

# 保存记忆
memory_id = str(uuid.uuid4())
memory = {"user": "我叫张三"}
store.put(namespace, memory_id, memory)

# 检索记忆
memories = store.search(namespace)
print(memories[-1].value)
# 输出:{'user': '我叫张三'}

6.3 结合 Checkpointer 和 Store

from langgraph.graph import StateGraph, MessagesState, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.store.memory import InMemoryStore
from langgraph.store.base import BaseStore
from langchain_core.runnables import RunnableConfig
import uuid

# 创建 checkpointer 和 store
memory = MemorySaver()
store = InMemoryStore()

# 定义节点(可以访问 store)
def call_model(
    state: MessagesState,
    config: RunnableConfig,
    *,
    store: BaseStore  # 注入 store
):
    # 获取用户 ID
    user_id = config["configurable"]["user_id"]
    namespace = ("memories", user_id)

    # 检索该用户的所有记忆
    memories = store.search(namespace)
    info = "\n".join([d.value["data"] for d in memories])

    # 使用记忆增强提示
    system_msg = f"用户的历史信息:{info}"

    # 调用 LLM
    response = llm.invoke(
        [{"type": "system", "content": system_msg}] + state["messages"]
    )

    # 保存新记忆
    last_message = state["messages"][-1]
    store.put(namespace, str(uuid.uuid4()), {"data": last_message.content})
    store.put(namespace, str(uuid.uuid4()), {"data": response.content})

    return {"messages": [response]}

# 构建图
builder = StateGraph(MessagesState)
builder.add_node("call_model", call_model)
builder.add_edge(START, "call_model")
builder.add_edge("call_model", END)

# 编译(传入 checkpointer 和 store)
graph = builder.compile(checkpointer=memory, store=store)
测试跨线程记忆
# 线程1:用户介绍自己
config1 = {
    "configurable": {
        "thread_id": "thread_1",
        "user_id": "user_123"  # 关键:user_id
    }
}
graph.invoke({"messages": [HumanMessage("我叫张三")]}, config1)

# 线程2:同一用户,不同线程
config2 = {
    "configurable": {
        "thread_id": "thread_2",
        "user_id": "user_123"  # 相同的 user_id
    }
}
result = graph.invoke({"messages": [HumanMessage("我叫什么?")]}, config2)
print(result["messages"][-1].content)
# 输出:你之前说你叫张三。(跨线程记忆成功!)
查看存储的记忆
# 查看用户的所有记忆
for memory in store.search(("memories", "user_123")):
    print(memory.value)

# 输出:
# {'data': '我叫张三'}
# {'data': '你好,张三!'}
# {'data': '我叫什么?'}
# {'data': '你之前说你叫张三。'}
实践技巧

Q1:MemorySaver 和 SqliteSaver 有什么区别?

A:

特性 MemorySaver SqliteSaver
存储位置 内存 数据库文件
持久化 ❌ 程序重启后丢失 ✅ 永久保存
适用场景 演示、测试 生产环境
性能 更快 稍慢(涉及磁盘 I/O)

Q2:如何清空对话历史?

A:使用新的 thread_id 或手动删除数据:

# 方法1:使用新 thread_id
config_new = {"configurable": {"thread_id": "new_conversation"}}

# 方法2:删除 SQLite 数据库文件
import os
os.remove("chatbot.db")

Q3:一个用户可以有多个 thread 吗?

A:可以!这是常见的设计模式:

# 用户 A 的不同会话
config_a1 = {"configurable": {"thread_id": "user_a_session_1"}}
config_a2 = {"configurable": {"thread_id": "user_a_session_2"}}

# 使用 Store 跨会话共享用户信息
user_id = "user_a"
namespace = ("users", user_id)
store.put(namespace, "name", {"value": "张三"})

Q4:如何限制对话历史长度?

A:在节点中手动截断消息:

def agent_with_trim(state: MessagesState):
    messages = state["messages"]

    # 只保留最近 10 条消息
    if len(messages) > 10:
        messages = messages[-10:]

    response = llm.invoke(messages)
    return {"messages": [response]}

Q5:Checkpointer 会影响性能吗?

A:会有一定影响,但通常可以接受:

  • MemorySaver:性能影响很小
  • SqliteSaver:涉及磁盘 I/O,稍慢
  • 优化建议:使用异步版本 AsyncSqliteSaver

 


人机交互

完全自主 Agent 的风险

我们已经学会构建完全自主的 Agent,它可以:

  • 调用外部工具
  • 维护多轮对话
  • 自主循环决策

但这带来一个严重问题:Agent 可能做出不当的决策

危险场景:

用户:"删除生产环境数据库"
Agent: [自动执行删除] ❌ 灾难!

用户:"转账10万元到..."
Agent: [自动执行转账] ❌ 资金损失!

用户:"修改系统配置"
Agent: [自动执行修改] ❌ 服务中断!

人机交互的价值

Human-in-the-loop (HIL) 允许在关键操作前引入人工审批:

用户请求 → Agent 判断 → [需要审批] → 暂停 → 人工确认 → 继续执行

价值:

  • ✅ 保持 Agent 自主性
  • ✅ 避免高风险操作
  • ✅ 增强系统可控性
  • ✅ 提升用户信任

适用场景:

  • 删除数据库操作
  • 金融交易确认
  • 机票改签审批
  • 配置变更确认
  • 敏感信息访问

类比理解:

就像核武器发射需要"双人确认"机制,重要的 AI 操作也需要人工"把关"。

断点(Breakpoint)机制

Breakpoint(断点): 在特定节点暂停图的执行,等待人工介入。

工作原理:

节点A → 节点B(断点)[暂停] → 人工操作 → [恢复] → 节点C

关键组件:

  1. Checkpointer: 保存图的状态(必需!)
  2. interrupt_before: 在节点执行前暂停
  3. interrupt_after: 在节点执行后暂停
  4. update_state(): 更新状态后恢复执行

断点类型

类型 设置方式 暂停时机 适用场景
静态断点 interrupt_before=["node"] 固定节点前 已知的审批点
动态断点 节点内部逻辑判断 条件触发时 灵活的审批逻辑

❓ 问题1:Breakpoint模块是干什么用的?

Breakpoint 是 LangGraph 实现 HIL 的核心机制。

核心功能:

  1. 暂停执行:在关键节点停止图的运行
  2. 保留状态:通过 checkpointer 保存当前状态
  3. 允许修改:人工可以查看和修改状态
  4. 恢复执行:修改后从暂停处继续

核心价值:

在保持 Agent 自主性的同时,引入人工监督和控制

❓ 问题2:最常用的方法是什么?

1. interrupt_before:节点前暂停
from langgraph.checkpoint.memory import MemorySaver

# 创建 checkpointer(必需!)
memory = MemorySaver()

# 在 execute_users 节点前暂停
graph = builder.compile(
    checkpointer=memory,
    interrupt_before=["execute_users"]  # 关键!
)

使用场景:

  • 在执行敏感操作前审批
  • 需要人工提供额外信息
  • 预防性质的暂停
2. interrupt_after:节点后暂停
graph = builder.compile(
    checkpointer=memory,
    interrupt_after=["agent"]  # agent 执行后暂停
)

使用场景:

  • 审查 Agent 的输出
  • 确认生成的内容
  • 验证工具调用参数
3. 运行到断点
config = {"configurable": {"thread_id": "1"}}

# 运行图,会在断点处停止
for chunk in graph.stream({"messages": [HumanMessage("删除数据")]}, config):
    print(chunk)

# 输出会在断点前停止
4. 查看状态
# 获取当前状态
snapshot = graph.get_state(config)

print(snapshot.values)      # 当前状态值
print(snapshot.next)        # 下一步要执行的节点
print(snapshot.tasks)       # 待执行的任务
5. 更新状态并恢复

方式1:直接恢复(同意操作)

# 直接恢复,不修改状态
for chunk in graph.stream(None, config):  # None = 继续
    print(chunk)

❓ 问题3:参数是什么意思?

compile() 参数
graph.compile(
    checkpointer: BaseCheckpointSaver,  # 必需!保存状态
    interrupt_before: list[str] = None, # 在哪些节点前暂停
    interrupt_after: list[str] = None   # 在哪些节点后暂停
)

interrupt_before/after:

  • 接收节点名称列表
  • 可以指定多个节点:["node1", "node2"]
  • 特殊值:["*"] 表示所有节点
update_state() 参数
graph.update_state(
    config: dict,              # 配置(包含 thread_id)
    values: dict,              # 要更新的状态值
    as_node: str = None        # 作为哪个节点更新
)

as_node 参数:

  • None(默认):直接更新状态
  • "node_name":模拟该节点的输出
  • 常用于注入工具响应

返回值:

{
    'configurable': {
        'thread_id': '1',
        'checkpoint_id': 'xxx...'  # 新的检查点 ID
    }
}

实战案例:敏感操作审批系统

构建一个数据管理系统,能够:

  1. 正常查询操作自动执行
  2. 删除操作需要人工审批
  3. 审批通过后才执行删除
  4. 审批拒绝后返回提示

4.2 完整实现

Step 1:定义状态
from typing import TypedDict

class State(TypedDict):
    user_input: str          # 用户输入
    model_response: str      # 模型响应
    user_approval: str       # 审批状态
Step 2:定义节点
from langchain_openai import ChatOpenAI
from langchain_core.messages import AIMessage

llm = ChatOpenAI(model="gpt-4")

def call_model(state):
    """处理用户请求"""
    messages = state["user_input"]

    # 判断是否是删除操作
    if '删除' in state["user_input"]:
        # 需要人工审批
        state["user_approval"] = f"用户输入的指令是:{state['user_input']}, 请人工确认是否执行!"
    else:
        # 普通操作,直接执行
        response = llm.invoke(messages)
        state["user_approval"] = "直接运行!"
        state["model_response"] = response

    return state

def execute_users(state):
    """执行用户操作(根据审批结果)"""
    if state["user_approval"] == "是":
        response = "您的删除请求已经获得管理员的批准并成功执行。"
        return {"model_response": AIMessage(response)}
    elif state["user_approval"] == "否":
        response = "对不起,您当前的请求是高风险操作,管理员不允许执行!"
        return {"model_response": AIMessage(response)}
    else:
        # 普通操作,直接返回
        return state
Step 3:构建图
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

# 构建图
builder = StateGraph(State)
builder.add_node("call_model", call_model)
builder.add_node("execute_users", execute_users)

builder.add_edge(START, "call_model")
builder.add_edge("call_model", "execute_users")
builder.add_edge("execute_users", END)

# 编译(添加断点!)
memory = MemorySaver()
graph = builder.compile(
    checkpointer=memory,
    interrupt_before=["execute_users"]  # 在执行前暂停
)
Step 4:测试删除操作
config = {"configurable": {"thread_id": "1"}}

# 运行到断点
for chunk in graph.stream(
    {"user_input": "我将在数据库中删除 id 为 xigualaoshi 的所有信息"},
    config
):
    print(chunk)

# 输出:
# {'user_input': '我将在数据库中删除...',
#  'user_approval': '用户输入的指令是:...请人工确认是否执行!'}
# [暂停在 execute_users 前]
Step 5:人工审批

情况1:同意删除

# 获取状态
snapshot = graph.get_state(config)

# 修改为同意
snapshot.values['user_approval'] = '是'
graph.update_state(config, snapshot.values)

# 恢复执行
for chunk in graph.stream(None, config):
    print(chunk)

# 输出:您的删除请求已经获得管理员的批准并成功执行。

情况2:拒绝删除

# 修改为拒绝
snapshot.values['user_approval'] = '否'
graph.update_state(config, snapshot.values)

# 恢复执行
for chunk in graph.stream(None, config):
    print(chunk)

# 输出:对不起,您当前的请求是高风险操作,管理员不允许执行!
Step 6:测试普通操作
# 普通问答,不触发断点
for chunk in graph.stream(
    {"user_input": "你好,请介绍一下你自己"},
    config
):
    print(chunk)

# 直接恢复(因为 user_approval = "直接运行!")
for chunk in graph.stream(None, config):
    print(chunk)

# 输出:正常的介绍内容

 


动态断点:工具调用场景

需求:一个 Agent 有多个工具,只对"删除"工具需要审批。

问题:

  • 使用 interrupt_before=["tools"] 会导致所有工具都暂停 ❌
  • 需要动态决定是否暂停 ✅

5.2 解决方案:条件路由 + 专用节点

架构设计
用户输入 → agent(判断) →
    ├─ 普通工具 → action → agent
    └─ 删除工具 → [断点] → run_tool → agent
完整实现

Step 1:定义工具

from langchain_core.tools import tool

@tool
def get_weather(city: str) -> str:
    """查询天气(安全操作)"""
    return f"{city}的天气:晴天"

@tool
def delete_data(city: str) -> str:
    """删除数据(危险操作!)"""
    return f"已删除 {city} 的数据"

tools = [get_weather, delete_data]

Step 2:定义路由函数

from langgraph.prebuilt import ToolNode

tool_node = ToolNode(tools)

def should_continue(state):
    """根据工具调用决定路由"""
    messages = state["messages"]
    last_message = messages[-1]

    if not last_message.tool_calls:
        return "end"

    # 检查是否调用删除工具
    if last_message.tool_calls[0]["name"] == "delete_data":
        return "run_tool"  # 需要审批
    else:
        return "continue"  # 直接执行

Step 3:定义专用工具节点

def run_tool(state):
    """执行删除工具(仅在此节点中执行)"""
    new_messages = []
    tool_calls = state["messages"][-1].tool_calls

    # 只处理删除工具
    tools_dict = {"delete_data": delete_data}

    for tool_call in tool_calls:
        tool = tools_dict[tool_call["name"]]
        result = tool.invoke(tool_call["args"])
        new_messages.append({
            "role": "tool",
            "name": tool_call["name"],
            "content": result,
            "tool_call_id": tool_call["id"]
        })

    return {"messages": new_messages}

Step 4:构建图

from langgraph.graph import MessagesState

workflow = StateGraph(MessagesState)

workflow.add_node("agent", call_model)
workflow.add_node("action", tool_node)      # 普通工具
workflow.add_node("run_tool", run_tool)     # 删除工具

workflow.add_edge(START, "agent")

workflow.add_conditional_edges(
    "agent",
    should_continue,
    {
        "continue": "action",      # 普通工具
        "run_tool": "run_tool",    # 删除工具
        "end": END
    }
)

workflow.add_edge("action", "agent")
workflow.add_edge("run_tool", "agent")

# 只在 run_tool 前设置断点!
memory = MemorySaver()
graph = workflow.compile(
    checkpointer=memory,
    interrupt_before=["run_tool"]  # 动态断点
)

Step 5:测试

测试1:普通工具(不暂停)

config = {"configurable": {"thread_id": "1"}}

for chunk in graph.stream(
    {"messages": [HumanMessage("查询北京天气")]},
    config
):
    chunk["messages"][-1].pretty_print()

# 输出:
# [调用 get_weather 工具]
# [直接返回结果,无暂停]

测试2:删除工具(暂停)

for chunk in graph.stream(
    {"messages": [HumanMessage("删除北京数据")]},
    config
):
    chunk["messages"][-1].pretty_print()

# 输出:
# [调用 delete_data 工具]
# [在 run_tool 前暂停]

# 同意删除
for chunk in graph.stream(None, config):
    chunk["messages"][-1].pretty_print()

# 输出:已删除北京的数据

测试3:拒绝删除(注入响应)

# 拒绝删除
state = graph.get_state(config)
tool_call_id = state.values["messages"][-1].tool_calls[0]["id"]

# 注入拒绝消息
new_message = {
    "role": "tool",
    "content": "管理员不允许执行该操作!",
    "name": "delete_data",
    "tool_call_id": tool_call_id
}

graph.update_state(config, {"messages": [new_message]}, as_node="run_tool")

# 恢复执行
for chunk in graph.stream(None, config):
    chunk["messages"][-1].pretty_print()

# 输出:管理员不允许执行该操作!

常见问题

Q1:忘记添加 checkpointer 会怎样?

A:会报错!断点必须依赖 checkpointer。

# ❌ 错误
graph = builder.compile(interrupt_before=["node"])

# ValueError: Checkpointer requires...

# ✅ 正确
memory = MemorySaver()
graph = builder.compile(checkpointer=memory, interrupt_before=["node"])

Q2:如何知道图在哪个节点暂停?

A:查看 snapshot.next:

snapshot = graph.get_state(config)
print(snapshot.next)  # ('execute_users',)

Q3:可以修改多个字段吗?

A:可以!update_state 可以修改任意字段:

graph.update_state(config, {
    "user_approval": "是",
    "user_name": "张三",
    "timestamp": "2025-01-01"
})

Q4:interrupt_before 和 interrupt_after 有什么区别?

A:

类型 暂停时机 适用场景
interrupt_before 节点执行前 预防性审批
interrupt_after 节点执行后 审查输出结果
# 在执行前审批
interrupt_before=["delete"]  # 删除前确认

# 在执行后审查
interrupt_after=["agent"]    # 查看 Agent 输出后决定

Q5:如何实现超时自动拒绝?

A:可以在外部循环中添加超时逻辑:

import time

# 记录暂停时间
start_time = time.time()

while True:
    state = graph.get_state(config)

    # 检查是否超时(例如 60 秒)
    if time.time() - start_time > 60:
        # 超时,自动拒绝
        tool_call_id = state.values["messages"][-1].tool_calls[0]["id"]
        new_message = {
            "role": "tool",
            "content": "操作超时,已自动拒绝",
            "tool_call_id": tool_call_id
        }
        graph.update_state(config, {"messages": [new_message]}, as_node="tools")
        break

    # 继续等待用户输入...

 


规划逻辑

在前面的章节中,我们构建的图都是线性流程:

START → 节点A → 节点B → 节点C → END

问题:

  • ❌ 无法根据不同输入选择不同路径
  • ❌ 所有请求都走相同流程
  • ❌ 无法实现智能决策
  • ❌ 灵活性差

1.2 规划逻辑的价值

通过添加条件路由,实现智能决策:

用户输入 → 判断 →
    ├─ 查询数据 → 数据库节点 → END
    ├─ 普通对话 → 聊天节点 → END
    └─ 调用工具 → 工具节点 → 聊天节点 → END

价值:

  • ✅ 根据输入智能选择路径
  • ✅ 不同任务走不同流程
  • ✅ 提升系统灵活性
  • ✅ 支持复杂的业务逻辑

适用场景:

  • 智能客服路由(售前/售后/投诉)
  • 多工具Agent(搜索/计算/数据库)
  • 审批流程(审批通过/拒绝/转交)
  • 任务分流(简单任务/复杂任务)

类比理解:

就像导航系统会根据路况选择不同路线,规划逻辑让Agent可以根据输入智能选择执行路径。

Conditional Edge(条件边): 根据节点输出动态决定下一个执行的节点

工作原理:

节点A → [Router Function判断] →
    ├─ 条件1为真 → 节点B
    ├─ 条件2为真 → 节点C
    └─ 其他情况 → 节点D

关键组件:

  1. Router Function: 判断函数,返回下一个节点名称
  2. Path Map: 将判断结果映射到节点名称
  3. Source Node: 条件边的起始节点

2.2 普通边 vs 条件边

类型 特点 使用场景
普通边 固定路径,总是到同一节点 确定的流程
条件边 动态路径,根据条件选择节点 需要决策的流程
# 普通边
builder.add_edge("node_a", "node_b")  # 总是 A → B

# 条件边
builder.add_conditional_edges(
    "node_a",
    router_function,
    {True: "node_b", False: "node_c"}
)  # 根据条件选择 B 或 C

❓ 问题1:add_conditional_edges这个方法是干什么用的?

add_conditional_edges 是 LangGraph 实现条件路由的核心方法。

核心功能:

  1. 动态路由:根据节点输出选择下一步
  2. 智能决策:支持复杂的判断逻辑
  3. 灵活控制:实现分支流程

核心价值:

将简单的线性图转变为智能的决策系统

❓ 问题2:最常用的方法是什么?

1. 基本用法
from langgraph.graph import StateGraph, START, END

def router_function(state):
    """路由函数:根据状态决定下一步"""
    if state["x"] == 10:
        return "node_b"
    else:
        return "node_c"

builder = StateGraph(dict)
builder.add_node("node_a", lambda s: {"x": s["x"] + 1})
builder.add_node("node_b", lambda s: {"x": s["x"] - 2})
builder.add_node("node_c", lambda s: {"x": s["x"] + 5})

# 添加条件边
builder.add_conditional_edges(
    "node_a",           # 起始节点
    router_function,    # 路由函数
)

builder.add_edge(START, "node_a")
builder.add_edge("node_b", END)
builder.add_edge("node_c", END)

graph = builder.compile()
2. 使用路径映射
def router_function(state):
    """返回布尔值"""
    return state["x"] > 10

builder.add_conditional_edges(
    "node_a",
    router_function,
    {
        True: "node_b",   # 当函数返回True时,去node_b
        False: "node_c"   # 当函数返回False时,去node_c
    }
)
3. 支持多个分支
def router_function(state):
    """返回分类标签"""
    if "查询" in state["user_input"]:
        return "search"
    elif "计算" in state["user_input"]:
        return "calculate"
    else:
        return "chat"

builder.add_conditional_edges(
    "classifier",
    router_function,
    {
        "search": "search_node",
        "calculate": "calc_node",
        "chat": "chat_node"
    }
)

❓ 问题3:参数是什么意思?

add_conditional_edges() 参数
builder.add_conditional_edges(
    source: str,                  # 起始节点名称
    path: Callable,               # 路由函数
    path_map: Optional[Dict] = None,  # 路径映射字典
    then: Optional[str] = None    # 可选:所有路径后的统一节点
)

source:

  • 条件边的起始节点
  • 必须是已添加的节点名称

path (路由函数):

  • 接收当前状态作为参数
  • 返回值决定下一个节点
  • 可以返回节点名称或任意可哈希值
def my_router(state) -> str:
    # 可以是复杂的判断逻辑
    if condition_1:
        return "node_a"
    elif condition_2:
        return "node_b"
    else:
        return "node_c"

path_map (可选):

  • 将路由函数返回值映射到节点名称
  • 如果不提供,路由函数必须直接返回节点名称
# 不使用path_map
def router(state):
    return "my_node"  # 直接返回节点名

# 使用path_map
def router(state):
    return True  # 返回任意值

path_map = {True: "node_a", False: "node_b"}

then (可选):

  • 所有条件路径执行后统一跳转的节点
  • 用于汇聚多条路径
builder.add_conditional_edges(
    "classifier",
    router_function,
    {
        "path_a": "node_a",
        "path_b": "node_b"
    },
    then="merge_node"  # node_a和node_b都会跳到这里
)

 


多 Agent

在前面的章节中,我们学会了构建强大的ReAct Agent,它可以:

  • 调用多个工具
  • 维护对话记忆
  • 自主循环决策

但随着任务复杂度增加,Single-Agent面临严峻挑战:

问题:

  • 工具过多导致混乱: Agent有太多工具可选,难以做出正确决策
  • 上下文过于复杂: 单一Agent需要处理所有信息,容易丢失关键细节
  • 缺乏专业化: 无法为不同领域设置专门的角色和背景

危险场景:

用户: "帮我分析这个数据集,生成报告,然后发送邮件"

Single-Agent:
- 需要理解数据分析工具
- 需要理解报告生成工具
- 需要理解邮件发送工具
- 容易在复杂流程中出错 ❌

让一个 Agent 在长时任务中持续、稳定地推进,远比想象困难。上下文窗口有限、会话之间没有记忆、反复返工、误判进度等问题,会让一个复杂项目在数轮执行后彻底失控。

Anthropic 发布了一套非常重要的工程方案,专门针对这些挑战而设计:基于“Initializer Agent + Coding Agent”的双 Agent 架构。

它的意义在于,它不是通过更大的模型、更长的上下文来对抗问题,而是通过一种工程化的工作流设计,确保 Agent 即使在“失忆”的多窗口条件下,也能像人类工程师一样一步步推进任务。

单 Agent 架构无法胜任长时任务,问题:

  • 一口气完成所有事情
  • 模型在看到部分成果后,误以为项目已经完成。由于缺乏清晰的目标列表与结构化任务定义,得到功能齐全的误解

Anthropic认为,当前AI Agent无法长时间稳定运行的核心问题不在模型能力,而在缺乏一种能够跨上下文继承任务逻辑的结构化工作方式。

双 Agent 架构:

  • 让 Agent 真正“像一个工程团队工作”Initializer Agent:一次性奠定整个项目的工程基础
  • Coding Agent:每一轮只做一件事,但把它做好

Anthropic 的解决方案非常工程化:不是让一个 Agent 解决所有事情,而是将长时任务拆分为两种角色——一个负责奠基,一个负责迭代。

Initializer Agent 的职责集中在第一次运行,它更像是一位“首席架构师”。

它不会立即进入编码,而是根据用户给出的高层需求,将项目转换为一个可长期维护的工程结构。

  • 第一,生成一个“可操作的需求体系”。

将用户需求分解为一份详尽的 JSON 结构的功能清单。每项功能都有描述、步骤和验收条件,并全部标记为“未完成”。

  • 第二,创建状态记录机制。

写入一个 progress 文件,记录项目结构、重要说明和之后用于交接的上下文。它同时建立 git 仓库,让每个迭代都能被提交、恢复和追踪。状态记录机制使得未来的 Agent 不必猜测,而是能够基于事实继续推进工作。

  • 第三,提供一个标准化的启动脚本。

Initializer 还会生成一个 init.sh,用于启动开发服务器并进行基础测试。这使得后续所有会话都能迅速验证当前环境是否健康,从而降低“接手时发现项目已坏”的风险。

Subgraphs:让Agent成为节点

在前面章节中,我们知道:

  • 节点(Node): 图中的基本单元,执行具体操作
  • 图(Graph): 由多个节点组成的工作流

Subgraphs(子图): 把一个完整的图作为另一个图中的节点使用。

父图:
  ├─ 简单节点A
  ├─ 简单节点B
  └─ 子图C (本身是一个完整的Agent图)
       ├─ 子图内部节点1
       ├─ 子图内部节点2
       └─ 子图内部节点3

关键问题: 父图和子图如何通信?

2.2 状态通信的两种方式

方式1: 共享状态键(推荐)

父图和子图有共同的状态字段:

# 父图状态
class ParentState(TypedDict):
    shared_data: str      # 共享字段
    parent_only: str      # 父图独有

# 子图状态
class SubgraphState(TypedDict):
    shared_data: str      # 共享字段(与父图相同)
    subgraph_only: str    # 子图独有

工作原理:

  • 父图通过shared_data传递数据给子图
  • 子图通过shared_data返回结果给父图
  • 各自的独有字段互不影响

示例:

from langgraph.graph import StateGraph, START, END
from typing import TypedDict

# 1. 定义状态
class ParentState(TypedDict):
    user_input: str
    final_answer: str

class SubgraphState(TypedDict):
    final_answer: str      # 共享键
    summary_answer: str    # 子图独有

# 2. 定义节点
def parent_node(state: ParentState):
    response = llm.invoke(state["user_input"])
    return {"final_answer": response}

def subgraph_node_1(state: SubgraphState):
    # 读取父图传来的数据
    full_answer = state["final_answer"]

    # 生成摘要
    summary = llm.invoke(f"总结以下内容:{full_answer}")
    return {"summary_answer": summary}

def subgraph_node_2(state: SubgraphState):
    # 评分
    score = llm.invoke(f"给这个回答打分:{state['summary_answer']}")

    # 更新共享字段,返回给父图
    return {"final_answer": score}

# 3. 构建子图
subgraph_builder = StateGraph(SubgraphState)
subgraph_builder.add_node("summary", subgraph_node_1)
subgraph_builder.add_node("score", subgraph_node_2)
subgraph_builder.add_edge(START, "summary")
subgraph_builder.add_edge("summary", "score")
subgraph = subgraph_builder.compile()

# 4. 构建父图(将子图作为节点)
builder = StateGraph(ParentState)
builder.add_node("parent", parent_node)
builder.add_node("subgraph", subgraph)  # 关键:子图作为节点
builder.add_edge(START, "parent")
builder.add_edge("parent", "subgraph")
builder.add_edge("subgraph", END)

graph = builder.compile()

# 5. 运行
result = graph.invoke({"user_input": "什么是机器学习?"})
print(result["final_answer"])  # 输出:子图返回的评分
方式2: 状态转换函数

父图和子图没有共同字段,需要手动转换:

def parent_node_with_transform(state: ParentState):
    # 1. 将父图状态转换为子图状态
    subgraph_input = {"different_key": state["parent_key"]}

    # 2. 调用子图
    subgraph_result = subgraph.invoke(subgraph_input)

    # 3. 将子图结果转换回父图状态
    return {"parent_key": subgraph_result["different_key"]}

# 父图中使用转换函数作为节点
builder.add_node("transform_node", parent_node_with_transform)

对比:

方式 优点 缺点
共享状态键 简单直接,自动传递 需要统一字段名称
状态转换 灵活,状态完全独立 需要手动编写转换逻辑

推荐: 优先使用共享状态键,只在状态结构差异大时使用转换函数。


多Agent架构模式

LangGraph提供了多种多Agent协作模式:

Network(网络)架构

特点: 每个Agent都可以与其他Agent通信,任何Agent都可以决定下一步调用谁。

    Agent A ←→ Agent B
       ↕          ↕
    Agent C ←→ Agent D

适用场景:

  • Agent数量较少(3-5个)
  • 需要灵活的通信模式
  • 任务流程不固定

优点:

  • ✅ 最大灵活性
  • ✅ Agent可以自主决策

缺点:

  • ❌ 难以扩展(Agent多时难以管理)
  • ❌ 可能产生无限循环
  • ❌ 难以追踪执行流程

Supervisor(主管)架构

特点: 一个Supervisor Agent负责协调多个Worker Agent,所有通信都经过Supervisor。

        Supervisor
       ↙    ↓    ↘
  Agent A  Agent B  Agent C

适用场景:

  • Agent数量较多
  • 需要集中控制
  • 有明确的任务分配逻辑

优点:

  • ✅ 清晰的控制流程
  • ✅ 易于扩展(添加新Worker)
  • ✅ 避免混乱的通信

缺点:

  • ❌ Supervisor成为瓶颈
  • ❌ 增加额外的决策开销

Hierarchical(分层)架构

特点: Supervisor嵌套Supervisor,形成多层管理结构。

     Top Supervisor
    ↙              ↘
Supervisor A    Supervisor B
↙  ↓  ↘        ↙  ↓  ↘
A1 A2 A3      B1 B2 B3

适用场景:

  • 大规模Agent系统
  • 复杂的业务层级
  • 需要模块化管理

推荐选择:

  • 3-5个Agent: Network架构
  • 5-10个Agent: Supervisor架构
  • 10+个Agent: Hierarchical架构

Network vs Supervisor对比

多Agent架构模式

  • Network: 灵活的多对多通信
  • Supervisor: 集中式调度管理
  • Hierarchical: 分层管理结构
维度 Network Supervisor
控制方式 分布式,Agent自主决策 集中式,Supervisor决策
通信模式 Agent之间直接通信 必须经过Supervisor
适用规模 3-5个Agent 5-10个Agent
灵活性 高,Agent可自由选择 中,由Supervisor控制
可维护性 低,通信关系复杂 高,只需维护Supervisor逻辑
性能开销 低,直接通信 中,多一层Supervisor决策
适用场景 需要灵活协作的小团队 需要明确分工的大团队

选择建议:

  • 如果Agent之间需要频繁灵活通信 → Network
  • 如果有清晰的任务分配逻辑 → Supervisor
  • 如果系统规模大,层级复杂 → Hierarchical(Supervisor嵌套)

常见问题

Q1: 如何避免Agent之间的无限循环?

A: 设置最大递归深度:

config = {"recursion_limit": 20}
graph.invoke(input, config)

Q2: 子图可以访问父图的所有状态吗?

A: 不可以。只能访问共享的状态键,或通过转换函数传递。

# 子图只能访问shared_key
class ParentState(TypedDict):
    shared_key: str
    private_key: str

class SubgraphState(TypedDict):
    shared_key: str  # 只能访问这个

Q3: Supervisor如何避免一直调用同一个Agent?

A: 在Supervisor的系统提示中明确指示:

system_prompt = """
如果某个worker已经执行过但没有解决问题,
应该尝试其他worker或直接FINISH,
避免重复调用同一个worker。
"""

Q4: 如何在多Agent系统中保持对话历史?

A: 使用Annotated[..., operator.add]:

class AgentState(TypedDict):
    messages: Annotated[Sequence[BaseMessage], operator.add]
    # 新消息会自动追加到列表

 


Debug 和监控

为什么需要Debug和监控?

随着Agent系统变得越来越复杂,调试和监控变得至关重要:

简单聊天机器人:
用户输入 → LLM → 回复 ✅ (容易调试)

复杂多Agent系统:
用户输入 → Agent A → 工具1 → Agent B → 工具2 → Agent C → 回复 ❌ (难以追踪)

问题:

  • 不透明: 不知道Agent在哪一步出错
  • 难以复现: 错误可能只在特定情况下发生
  • 性能瓶颈: 不清楚哪个步骤最耗时
  • Token消耗: 不知道哪里消耗了大量Token
  • 逻辑错误: Agent的决策路径不符合预期

通过有效的调试和监控,可以:

调试阶段:

  • 追踪执行流: 看到每个节点的执行顺序
  • 检查状态变化: 了解状态在各节点间的传递
  • 定位错误: 快速找到出错的节点和原因
  • 验证逻辑: 确保路由决策符合预期

生产阶段:

  • 性能监控: 识别慢节点和瓶颈
  • 成本控制: 监控Token使用量
  • 质量保证: 追踪Agent的表现
  • 问题排查: 快速定位线上问题

类比理解:

就像开车需要仪表盘显示速度、油量、故障灯一样,复杂的Agent系统也需要"仪表盘"来监控运行状态。

stream_mode=“debug” 模式

核心功能: 输出图执行过程中的所有详细信息

基本用法
from langgraph.graph import StateGraph, MessagesState, START, END
from langchain_openai import ChatOpenAI

# 构建简单的图
llm = ChatOpenAI(model="gpt-4")

def chatbot(state: MessagesState):
    response = llm.invoke(state["messages"])
    return {"messages": [response]}

builder = StateGraph(MessagesState)
builder.add_node("chatbot", chatbot)
builder.add_edge(START, "chatbot")
builder.add_edge("chatbot", END)
graph = builder.compile()

# 使用debug模式
for chunk in graph.stream(
    {"messages": [HumanMessage("你好")]},
    stream_mode="debug"
):
    print(chunk)
Debug输出解析
# 输出示例
{
    'type': 'task',                    # 事件类型
    'timestamp': '2025-01-02T...',    # 时间戳
    'step': 1,                         # 步骤编号
    'payload': {
        'id': 'xxx...',                # 任务ID
        'name': 'chatbot',             # 节点名称
        'input': {...},                # 输入状态
        'triggers': ['start:chatbot']  # 触发来源
    }
}

{
    'type': 'task_result',             # 任务结果
    'step': 1,
    'payload': {
        'name': 'chatbot',
        'result': [...],               # 返回结果
        'error': None                  # 错误信息(如有)
    }
}

关键字段说明:

字段 说明 用途
type 事件类型(task/task_result) 区分开始和结束
step 超级步骤编号 追踪执行顺序
name 节点名称 定位具体节点
input 节点输入状态 检查传入数据
result 节点输出结果 验证返回值
error 错误信息 快速定位问题

LangSmith 是LangChain官方提供的监控和调试平台:

开发阶段                生产阶段
   ↓                       ↓
调试追踪 ←─── LangSmith ───→ 性能监控
测试评估                  质量分析

核心功能:

  • 📊 可视化追踪: 图形化展示执行流程
  • 🔍 详细日志: 记录所有输入输出
  • 📈 性能分析: Token使用、延迟统计
  • 🧪 测试评估: 对比不同版本表现
  • 🔔 告警通知: 异常情况实时提醒

事件流 astream_events 深度调试

核心概念: 访问图执行过程中的所有事件

async for event in graph.astream_events(
    input_data,
    version="v2"  # 必须指定版本
):
    # 处理事件
    pass

事件类型

事件类型 说明 适用场景
on_chain_start 节点开始执行 追踪执行流程
on_chain_end 节点执行完成 检查最终结果
on_chain_stream 节点流式输出 实时监控进度
on_chat_model_start LLM开始调用 记录模型参数
on_chat_model_stream LLM流式输出 获取Token流
on_chat_model_end LLM调用结束 统计Token使用
on_tool_start 工具开始执行 追踪工具调用
on_tool_end 工具执行完成 验证工具结果

实战:监控Token使用

from langchain_core.messages import HumanMessage

async def monitor_token_usage(input_msg):
    """监控每个节点的Token使用情况"""
    token_stats = {}

    async for event in agent.astream_events(
        {"messages": [HumanMessage(input_msg)]},
        version="v2"
    ):
        kind = event["event"]

        # 只关注LLM调用结束事件
        if kind == "on_chat_model_end":
            name = event["name"]
            data = event["data"]

            # 提取Token使用量
            if "output" in data:
                output = data["output"]
                if hasattr(output, "usage_metadata"):
                    usage = output.usage_metadata
                    token_stats[name] = {
                        "input_tokens": usage.get("input_tokens", 0),
                        "output_tokens": usage.get("output_tokens", 0),
                        "total_tokens": usage.get("total_tokens", 0)
                    }

    return token_stats

# 使用
import asyncio

stats = asyncio.run(monitor_token_usage("北京天气怎么样,并分析一下适合穿什么衣服"))

print("\n=== Token使用统计 ===")
for model_name, usage in stats.items():
    print(f"\n{model_name}:")
    print(f"  输入Token: {usage['input_tokens']}")
    print(f"  输出Token: {usage['output_tokens']}")
    print(f"  总计: {usage['total_tokens']}")

常见问题

6.1 场景1:Agent陷入循环

问题: Agent反复调用相同的工具

# 症状
for chunk in graph.stream(input, stream_mode="updates"):
    print(chunk)

# 输出:
# {'agent': ...}
# {'tools': ...}
# {'agent': ...}  # 又回到agent
# {'tools': ...}  # 又调用tools
# ... (无限循环)

诊断方法:

# 添加循环检测
def debug_with_loop_detection(input_data, max_steps=10):
    steps = []

    for i, chunk in enumerate(graph.stream(input_data, stream_mode="updates")):
        node_name = list(chunk.keys())[0]
        steps.append(node_name)

        print(f"Step {i+1}: {node_name}")

        # 检测循环
        if len(steps) > max_steps:
            print(f"\n❌ 警告: 超过{max_steps}步,可能陷入循环!")
            print(f"执行路径: {' → '.join(steps)}")
            break

    return steps

# 使用
path = debug_with_loop_detection({"messages": [HumanMessage("测试")]})

解决方案:

# 1. 设置最大递归深度
config = {"recursion_limit": 20}
graph.invoke(input, config)

# 2. 改进路由逻辑
def should_continue(state):
    # 添加循环检测
    if len(state["messages"]) > 10:
        return END  # 强制结束

    last_msg = state["messages"][-1]
    if last_msg.tool_calls:
        return "tools"
    else:
        return END

6.2 场景2:状态传递错误

问题: 节点接收到的状态不符合预期

# 症状: KeyError或数据丢失
def my_node(state):
    data = state["expected_key"]  # KeyError!
    return {"result": data}

诊断方法:

# 在每个节点添加状态检查
def debug_node(state):
    print(f"\n=== 节点输入状态 ===")
    print(f"状态键: {list(state.keys())}")
    for key, value in state.items():
        print(f"{key}: {type(value)} = {value}")

    # 正常逻辑
    result = llm.invoke(state["messages"])
    return {"messages": [result]}

# 或使用debug模式查看
for chunk in graph.stream(input, stream_mode="debug"):
    if chunk['type'] == 'task':
        print(f"节点 {chunk['payload']['name']} 的输入:")
        print(chunk['payload']['input'])

解决方案:

# 1. 使用TypedDict明确状态结构
from typing import TypedDict, Annotated
from langgraph.graph.message import add_messages

class MyState(TypedDict):
    messages: Annotated[list, add_messages]
    user_data: dict  # 明确需要的字段

# 2. 在节点中添加防御性检查
def safe_node(state: MyState):
    if "user_data" not in state:
        return {"user_data": {}}  # 提供默认值

    # 继续处理
    ...

6.3 场景3:性能瓶颈

问题: 某个节点执行特别慢

诊断方法:

import time

async def measure_performance(input_data):
    """测量每个节点的执行时间"""
    timings = {}
    current_step = None
    start_time = None

    async for event in graph.astream_events(input_data, version="v2"):
        kind = event["event"]
        name = event["name"]

        if kind == "on_chain_start":
            current_step = name
            start_time = time.time()

        elif kind == "on_chain_end" and current_step == name:
            elapsed = time.time() - start_time
            timings[name] = elapsed
            print(f"{name}: {elapsed:.2f}秒")

    return timings

# 使用
timings = asyncio.run(measure_performance({"messages": [HumanMessage("测试")]}))

# 找出最慢的节点
slowest = max(timings.items(), key=lambda x: x[1])
print(f"\n最慢节点: {slowest[0]} ({slowest[1]:.2f}秒)")

解决方案:

# 1. 优化LLM调用
llm_fast = ChatOpenAI(
    model="gpt-3.5-turbo",  # 使用更快的模型
    temperature=0,
    streaming=True  # 启用流式输出
)

# 2. 并行化独立操作
from langgraph.graph import StateGraph

# 如果node_a和node_b独立,可以并行执行
builder.add_node("node_a", node_a)
builder.add_node("node_b", node_b)
builder.add_edge(START, "node_a")
builder.add_edge(START, "node_b")  # 并行
Logo

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

更多推荐