文章总结: 文档详细介绍了LangGraph1.0的核心改进,包括中间件机制通过钩子函数实现工作流可控扩展,状态图设计用有向图和状态机范式取代线性Agent循环,以及类型安全与ContextAPI升级。文章深入讲解了状态管理、StateGraphAPI、并行执行与超步、流式执行模式、设计模式与最佳实践,以及生产级应用架构,为开发者提供了构建复杂Agent系统的完整指南。
综合评分: 93
文章分类: AI安全,应用安全,安全开发
LangGraph 1.0 完全指南:新概念、API、设计模式与最佳实践
原创
黄师傅
黄师傅的赛博dojo
2025年12月6日 09:45
上海
#
看了https://www.luochang.ink/dive-into-langgraph/ ,让pplx做了一个pro plus版本,主要是给各种vibe coding工具看的,用来做一些架构审查和优化的工作。
执行摘要
LangGraph 1.0(发布于 2025年10月18日)标志着 Agent 框架从无序混乱向系统化工程的重大转变。核心改进围绕三大支柱展开:
- 1. 中间件机制(Middleware):通过钩子函数实现工作流的可控扩展,彻底解决了上下文工程混乱的问题
- 2. 状态图设计(StateGraph):用有向图 + 状态机范式取代线性 Agent 循环,支持复杂工作流编排
- 3. 类型安全与运行时控制:从 v0.6 的配置地狱升级到优雅的 Context API + Reducer 机制
本指南重点覆盖这三大核心创新的具体实现、最佳实践和生产级模式。
第一部分:LangGraph 1.0 的三大核心改进
1.1 从 v0.6 到 v1.0:问题诊断
v0.6 的核心痛点:
# ❌ v0.6 的"配置地狱"
def node(state: State, config: RunnableConfig):
# 层层嵌套获取数据,容易出错
user_id = config.get("configurable", {}).get("user_id")
db_conn = config.get("configurable", {}).get("db_connection")
# 代码可读性差,维护成本高
问题根源:
- • 上下文管理缺乏系统化支持,全靠手写配置
- • 工具调用权限难以精确控制
- • 长对话的消息管理容易爆炸
- • 没有统一的扩展机制(预算控制、审计、安全过滤等)
1.2 v1.0 的解决方案三支柱
支柱1:中间件机制 (Middleware)
核心理念: 用类似 FastAPI 中间件的概念,在 Agent 执行流程的关键环节插入钩子函数。
执行生命周期:
User Input
↓
[before_model] 中间件 - 预处理输入
↓
[wrap_model_call] 中间件 - 包裹模型调用、修改参数
↓
LLM Model Call
↓
[wrap_tool_call] 中间件 - 拦截工具调用、权限控制
↓
Tool Execution
↓
[after_model] 中间件 - 输出验证、安全检查
↓
Response
内置中间件示例:
from langchain.agents import create_agent
from langchain.agents.middleware import (
PIIMiddleware,
SummarizationMiddleware,
HumanInTheLoopMiddleware
)
agent = create_agent(
model="claude-sonnet-4-5-20250929",
tools=[read_email, send_email],
middleware=[
# 隐私保护:脱敏邮箱地址
PIIMiddleware("email", strategy="redact"),
# 隐私保护:阻止电话号码
PIIMiddleware("phone_number", strategy="block"),
# 自动摘要:对话超过500 tokens自动压缩
SummarizationMiddleware(
model="claude-sonnet-4-5-20250929",
max_tokens_before_summary=500
),
# 人机交互:发送邮件前要求人工审批
HumanInTheLoopMiddleware(
interrupt_on={
"send_email": {
"allowed_decisions": ["approve", "edit", "reject"]
}
}
)
]
)
自定义中间件开发模式:
from dataclasses import dataclass
from typing import Callable
from langchain.agents.middleware import AgentMiddleware, ModelRequest
from langchain.agents.middleware.types import ModelResponse
from langchain_openai import ChatOpenAI
@dataclass
class Context:
user_expertise: str = "beginner" # "beginner" 或 "expert"
class ExpertiseBasedToolMiddleware(AgentMiddleware):
"""根据用户技术水平动态调整AI能力的中间件"""
def wrap_model_call(
self,
request: ModelRequest,
handler: Callable[[ModelRequest], ModelResponse]
) -> ModelResponse:
# 读取运行时上下文
user_level = request.runtime.context.user_expertise
if user_level == "expert":
# 专家用户:强模型 + 高级工具
request.model = ChatOpenAI(model="gpt-5")
request.tools = [advanced_search, data_analysis, ml_toolkit]
else:
# 初学者:轻量模型 + 基础工具
request.model = ChatOpenAI(model="gpt-5-nano")
request.tools = [simple_search, basic_calculator]
# 继续执行流程
return handler(request)
# 使用自定义中间件
agent = create_agent(
model="claude-sonnet-4-5-20250929",
tools=[simple_search, advanced_search, basic_calculator, data_analysis],
middleware=[ExpertiseBasedToolMiddleware()],
context_schema=Context
)
中间件的设计优势:
- • ✅ 代码模块化,功能边界清晰
- • ✅ 复用性高,可跨项目共享
- • ✅ 组合灵活,如搭积木般构建复杂功能
- • ✅ 测试友好,每个中间件可独立测试
- • ✅ 生产友好,常见需求都有模式支持
支柱2:状态图设计(StateGraph + 状态机范式)
核心理念: 用有向图的节点和边来表示 Agent 工作流,状态在节点间流动。
从 ReAct Agent 到 StateGraph 的进化:
# ❌ ReAct Agent 方式:无序循环
from langchain.agents import create_openai_tools_agent, AgentExecutor
agent = create_openai_tools_agent(llm, tools, prompt)
executor = AgentExecutor(agent=agent, tools=tools, max_iterations=10)
result = executor.invoke({"input": user_query})
# 问题:流程不可控,调试困难,扩展受限
# ✅ StateGraph 方式:可控的状态机
from langgraph.graph import StateGraph, MessagesState, START, END
# 定义状态(所有节点间的共享数据)
class State(TypedDict):
messages: Annotated[list, add_messages] # 消息列表,自动追加
user_id: str
context: dict
# 定义节点(处理单元)
def assistant_node(state: State) -> dict:
"""LLM 推理节点"""
response = llm.invoke(state["messages"])
return {"messages": [response]}
def tool_node(state: State) -> dict:
"""工具执行节点"""
last_msg = state["messages"][-1]
tool_results = execute_tools(last_msg.tool_calls)
return {"messages": [ToolMessage(content=tool_results)]}
# 路由函数:根据状态决定下一步
def route_after_llm(state: State) -> str:
last_msg = state["messages"][-1]
if last_msg.tool_calls:
return "tool"
return END
# 构建图
builder = StateGraph(State)
builder.add_node("assistant", assistant_node)
builder.add_node("tool", tool_node)
# 添加边
builder.add_edge(START, "assistant")
builder.add_conditional_edges("assistant", route_after_llm)
builder.add_edge("tool", "assistant")
# 编译为可执行的图
graph = builder.compile()
# 调用
result = graph.invoke({"messages": [HumanMessage("...")]})
优势对比:
| 特性 | ReAct Agent | StateGraph |
| — | — | — |
| 流程可视化 | ❌ 黑盒 | ✅ 有向图 |
| 并行执行 | ❌ 串行循环 | ✅ 超步并行 |
| 条件路由 | 有限 | ✅ 灵活丰富 |
| 状态管理 | 混乱 | ✅ 显式清晰 |
| 调试 | 困难 | ✅ 易调试 |
| 扩展性 | 差 | ✅ 优秀 |
支柱3:类型安全与上下文控制(Context API)
v0.6 的困境 → v1.0 的优雅:
# ❌ v0.6:配置参数层层嵌套
def node(state: State, config: RunnableConfig):
user_id = config.get("configurable", {}).get("user_id")
db_conn = config.get("configurable", {}).get("db_connection")
# 类型提示缺失,IDE 无法补全,运行时容易出错
# ✅ v1.0:Context API,类型安全
@dataclass
class Context:
user_id: str
db_connection: Connection
cache_client: Redis
def node(state: State, runtime: Runtime[Context]):
# 直接访问,IDE 自动补全,类型检查
user_id = runtime.context.user_id
db_conn = runtime.context.db_connection
第二部分:状态管理深度指南
2.1 State 的三种定义方式
方式1:普通字典(最灵活但易出错)
from langgraph.graph import StateGraph, START, END
def add_node(state):
# state 是普通字典,IDE 无法提示
return {"result": state["x"] + 1}
builder = StateGraph(dict) # 使用 dict 作为状态
builder.add_node("process", add_node)
方式2:TypedDict(推荐用于中小项目)
from typing import TypedDict, Annotated, List
from langgraph.graph import StateGraph
class State(TypedDict):
x: int # 单纯数值,默认覆盖
messages: List[str] # 列表字段
metadata: dict # 嵌套对象
def process(state: State) -> dict:
# IDE 能补全,类型检查
x = state["x"] # ✅ IDE 提示存在
# y = state["y"] # ❌ IDE 警告不存在
return {"x": x + 1}
builder = StateGraph(State)
方式3:Pydantic BaseModel(最严谨,用于大型项目)
from pydantic import BaseModel, Field
class State(BaseModel):
x: int = Field(description="计算数值")
messages: List[str] = Field(default_factory=list)
class Config:
# 允许额外字段
extra = "allow"
# 使用方式与 TypedDict 相同
builder = StateGraph(State)
2.2 Reducer 函数:状态融合的核心机制
问题场景: 多个节点同时修改同一字段,如何合并?
# ❌ 问题:不加 Reducer,后来的覆盖前面的
# 节点1 返回: {"messages": [msg1]}
# 节点2 返回: {"messages": [msg2]}
# 结果:messages 只有 [msg2],msg1 丢失!
# ✅ 解决:用 Reducer 指定合并策略
from typing import Annotated
import operator
from langgraph.graph.message import add_messages
class State(TypedDict):
# 注解语法:Annotated[类型, Reducer函数]
messages: Annotated[List, add_messages] # 智能合并
numbers: Annotated[List[int], operator.add] # 列表拼接
count: Annotated[int, lambda a, b: a + b] # 自定义:求和
常用 Reducer 函数详解:
import operator
from typing import Annotated, List
class State(TypedDict):
# 1. operator.add:列表或字符串拼接
messages: Annotated[List, operator.add]
# 2. add_messages:LangChain 专用,智能处理消息合并
# 特性:
# - 给消息分配 ID
# - ID 相同时更新,ID 不同时追加
# - 支持消息编辑(比如用户修改 AI 回复后继续对话)
from langgraph.graph.message import add_messages
messages: Annotated[List, add_messages]
# 3. 自定义 Reducer:选最大值
def max_reducer(a: int, b: int) -> int:
return max(a, b)
high_score: Annotated[int, max_reducer]
# 4. 自定义 Reducer:合并字典
def merge_dicts(a: dict, b: dict) -> dict:
result = a.copy()
result.update(b)
return result
config: Annotated[dict, merge_dicts]
# 实际应用示例
def node1(state: State) -> dict:
return {"messages": [msg1], "count": 1}
def node2(state: State) -> dict:
return {"messages": [msg2], "count": 5}
# 执行流程:
# 初始:{"messages": [], "count": 0}
# 节点1 后:{"messages": [msg1], "count": 1}
# 节点2 后:{"messages": [msg1, msg2], "count": 6} # operator.add 求和
add_messages 的智能特性:
from langgraph.graph.message import add_messages
from langchain_core.messages import HumanMessage, AIMessage
# 场景1:追加新消息
msg1 = [HumanMessage(content="你好", id="1")]
msg2 = [AIMessage(content="你好啊", id="2")]
result = add_messages(msg1, msg2)
# 结果:[msg1, msg2]
# 场景2:更新已有消息(ID 相同)
msg1 = [HumanMessage(content="你好", id="1")]
msg2 = [HumanMessage(content="你好呀", id="1")] # 相同ID
result = add_messages(msg1, msg2)
# 结果:[msg2] # 更新了第一条消息
# 场景3:删除消息(返回 None)
msg1 = [msg1, msg2]
msg3 = [{"id": "1", "type": "DELETE"}] # 删除 ID 为 1 的消息
result = add_messages(msg1, msg3)
# 结果:[msg2]
2.3 完整的状态定义最佳实践
from typing import TypedDict, Annotated, List, Optional
from dataclasses import dataclass
from langgraph.graph.message import add_messages
from langchain_core.messages import BaseMessage
import operator
# ✅ 最佳实践:完整的 State 定义
class ConversationState(TypedDict):
# 1. 消息管理:使用 add_messages 智能追加
messages: Annotated[List[BaseMessage], add_messages]
# 2. 用户上下文:保持不变
user_id: str
session_id: str
# 3. 工作流状态:记录处理过程
current_step: str # "收集需求" | "生成方案" | "反馈改进"
processing_status: str # "processing" | "completed" | "failed"
# 4. 中间结果:使用 Reducer 聚合
analysis_results: Annotated[List[dict], operator.add]
generated_solutions: Annotated[List[str], operator.add]
# 5. 计数器:使用求和 Reducer
total_tokens: Annotated[int, lambda a, b: a + b]
tool_calls_count: Annotated[int, lambda a, b: a + b]
# 6. 元数据:保留最新版本
metadata: dict # 每次返回时覆盖
# 7. 错误追踪:使用自定义 Reducer
def merge_errors(old: List[str], new: List[str]) -> List[str]:
"""合并错误列表,去重"""
return list(set(old + new))
errors: Annotated[List[str], merge_errors]
第三部分:StateGraph API 完全参考
3.1 核心方法详解
3.1.1 add_node:添加处理节点
from langgraph.graph import StateGraph, START, END
builder = StateGraph(State)
# 方式1:方法名作为节点名称
def process_data(state: State) -> dict:
"""处理函数名会自动变成节点名 'process_data'"""
return {"result": ...}
builder.add_node(process_data)
# 方式2:显式指定节点名称
builder.add_node("custom_name", process_data)
# 方式3:使用 Runnable 对象
from langchain_core.runnables import Runnable
runnable = some_chain.map()
builder.add_node("my_runnable", runnable)
# 返回值:更新后的 State 的子集(部分更新)
def example_node(state: State) -> dict:
# 不用返回整个 State,只返回更新的部分
return {
"messages": [...], # 会与原消息合并
"status": "done" # 覆盖原状态
}
3.1.2 add_edge:添加无条件边(节点连接)
# 语法:add_edge(source, target)
# source 可以是:字符串(节点名)、START、END
# target 可以是:字符串(节点名)、END
# 场景1:从起点到第一个节点
builder.add_edge(START, "step1")
# 场景2:从一个节点到另一个节点
builder.add_edge("step1", "step2")
# 场景3:节点到结束
builder.add_edge("final_step", END)
# 场景4:多个节点指向同一个下游节点(扇入)
builder.add_edge("parallel_task_a", "aggregator")
builder.add_edge("parallel_task_b", "aggregator")
builder.add_edge("parallel_task_c", "aggregator")
3.1.3 add_conditional_edges:条件边(动态路由)
def route_function(state: State) -> str:
"""
条件函数,返回下一个节点名称
可以实现:IF-ELSE、多分支路由、重试机制等
"""
if len(state["messages"]) < 3:
return "gather_more_info"
elif "error" in state["messages"][-1].content:
return "handle_error"
else:
return "generate_response"
# 添加条件边
builder.add_conditional_edges(
"decide_next_step", # 源节点
route_function, # 条件函数
{ # 目标节点映射
"gather_more_info": "info_collection",
"handle_error": "error_handler",
"generate_response": "response_gen"
}
)
# 高级用法:路由到 END
def route_with_end(state: State) -> str:
if state["processing_status"] == "completed":
return END # 直接结束流程
return "next_step"
builder.add_conditional_edges(
"check_completion",
route_with_end,
{
"next_step": "process",
END: END # 显式指向终点
}
)
3.1.4 add_sequence:序列快捷方式
# 等价写法对比
# ❌ 冗长方式
builder.add_edge(START, "step1")
builder.add_edge("step1", "step2")
builder.add_edge("step2", "step3")
builder.add_edge("step3", END)
# ✅ 简洁方式
builder.add_sequence(["step1", "step2", "step3"])
# 自动生成:START -> step1 -> step2 -> step3 -> END
3.1.5 compile:编译为可执行图
# 基本编译
compiled_graph = builder.compile()
# 带图名称和配置
compiled_graph = builder.compile(
name="my-workflow",
config={"recursion_limit": 25} # 防止无限循环
)
# 返回值:CompiledStateGraph 对象
# 包含方法:invoke(), stream(), astream(), get_graph() 等
3.2 特殊节点:START 和 END
from langgraph.graph import START, END
# START:图的入口点
# - 接收初始输入
# - 必须有至少一条出边
builder.add_edge(START, "first_node")
# END:图的出口点
# - 表示执行完成
# - 可以有多个节点指向它
builder.add_edge("step1", END)
builder.add_edge("step2", END) # 多个终点路由
第四部分:并行执行与超步(SuperStep)
4.1 超步概念
LangGraph 的 超步(SuperStep) 是并行执行的基本单位。
传统循环式 Agent:
step 1: 节点A (顺序执行)
step 2: 节点B
step 3: 节点C
LangGraph 超步:
超步1: 节点A (执行)
超步2: 节点B + 节点C (并行执行) ← 关键:无依赖则并行
超步3: 节点D (等待B、C完成)
4.2 实现并行执行的三种模式
模式1:多条输出边(扇出)
from langgraph.graph import StateGraph, START, END
class State(TypedDict):
data: str
results: Annotated[List[str], operator.add]
def split_task(state: State) -> str:
"""根据数据类型,决定并行任务"""
if "analyze" in state["data"]:
return "route_to_parallel"
return "route_to_sequential"
# 并行任务节点
def task_a(state: State) -> dict:
return {"results": ["Task A completed"]}
def task_b(state: State) -> dict:
return {"results": ["Task B completed"]}
def task_c(state: State) -> dict:
return {"results": ["Task C completed"]}
def aggregator(state: State) -> dict:
"""聚合所有并行任务的结果"""
print(f"汇总:{state['results']}")
return {"results": ["All tasks done"]}
# 构建图
builder = StateGraph(State)
# 添加节点
builder.add_node("split", split_task)
builder.add_node("task_a", task_a)
builder.add_node("task_b", task_b)
builder.add_node("task_c", task_c)
builder.add_node("aggregate", aggregator)
# 路由到并行任务
builder.add_edge(START, "split")
def route_to_parallel(state: State) -> str:
return "parallel"
builder.add_conditional_edges(
"split",
lambda s: "parallel",
{"parallel": "task_a"} # 实际需要复杂路由
)
# 关键:三个任务会在同一个超步中并行执行
builder.add_edge("split", "task_a")
builder.add_edge("split", "task_b")
builder.add_edge("split", "task_c")
# 扇入:等待所有并行任务完成后聚合
builder.add_edge("task_a", "aggregate")
builder.add_edge("task_b", "aggregate")
builder.add_edge("task_c", "aggregate")
builder.add_edge("aggregate", END)
graph = builder.compile()
# 执行
result = graph.invoke({"data": "analyze something", "results": []})
# 输出顺序:
# 超步1:split 执行
# 超步2:task_a, task_b, task_c 并行执行
# 超步3:aggregate 聚合
模式2:使用 Send 进行动态分发
from langgraph.graph import Send
class State(TypedDict):
tasks: List[str]
results: Annotated[List[str], operator.add]
def dispatcher(state: State) -> List[Send]:
"""
使用 Send 对象动态创建并行任务
返回 Send 对象列表,每个对应一个并行任务
"""
sends = []
for task_id, task in enumerate(state["tasks"]):
sends.append(
Send(
"process_task", # 目标节点
{"task": task, "task_id": task_id} # 传递给节点的状态
)
)
return sends
def process_task(state: State) -> dict:
"""处理单个任务"""
return {"results": [f"Processed: {state['task']}"]}
def collect_results(state: State) -> dict:
"""收集所有结果"""
return {"results": ["Collection complete"]}
# 构建图
builder = StateGraph(State)
builder.add_node("dispatcher", dispatcher)
builder.add_node("process_task", process_task)
builder.add_node("collector", collect_results)
builder.add_edge(START, "dispatcher")
builder.add_conditional_edges(
"dispatcher",
lambda x: x, # 返回 Send 对象列表
{"process_task": "process_task"} # 让框架处理路由
)
builder.add_edge("process_task", "collector")
builder.add_edge("collector", END)
graph = builder.compile()
# 执行
result = graph.invoke({
"tasks": ["task1", "task2", "task3"],
"results": []
})
# 三个任务会动态并行执行
模式3:defer 参数控制执行顺序
# 场景:需要等待所有并行任务完成后才执行聚合
builder = StateGraph(State)
# 添加节点
builder.add_node("task_a", task_a)
builder.add_node("task_b", task_b)
builder.add_node("task_c", task_c)
builder.add_node("aggregator", aggregator)
# 关键:设置 defer=True
# 这样 aggregator 会等待所有待处理任务完成
builder.add_node(
"aggregator",
aggregator,
defer=True # 延迟执行,等待所有依赖完成
)
# 构建并行流
builder.add_edge(START, "task_a")
builder.add_edge("task_a", "task_b")
builder.add_edge("task_b", "task_c")
builder.add_edge("task_c", "aggregator")
# 超步执行:
# 超步1:task_a
# 超步2:task_b(a 依赖完成)
# 超步3:task_c(b 依赖完成)
# 超步4:aggregator(defer 等待所有完成)
第五部分:流式执行与流模式(Stream Modes)
5.1 四种流模式对比
graph = builder.compile()
# 模式1:values - 完整状态流
print("=== values 模式 ===")
for chunk in graph.stream(
{"messages": [HumanMessage("Hi")]},
stream_mode="values"
):
print(chunk) # 每次输出完整的 State 字典
# 输出示例:
# {'messages': [msg1], 'status': 'processing'}
# {'messages': [msg1, msg2], 'status': 'done'}
# 模式2:updates - 增量更新
print("=== updates 模式 ===")
for chunk in graph.stream(
{"messages": [HumanMessage("Hi")]},
stream_mode="updates"
):
print(chunk) # 字典格式:{节点名: 更新值}
# 输出示例:
# {'assistant': {'messages': [msg1]}}
# {'tool': {'messages': [msg2]}}
# 模式3:debug - 详细调试信息
print("=== debug 模式 ===")
for chunk in graph.stream(
{"messages": [HumanMessage("Hi")]},
stream_mode="debug"
):
print(chunk) # 包含节点执行时间、错误、完整状态等
# 模式4:messages - LLM token 流(用于聊天应用)
print("=== messages 模式 ===")
for chunk in graph.stream(
{"messages": [HumanMessage("Hi")]},
stream_mode="messages"
):
print(chunk) # 只流式输出 LLM tokens
5.2 选择流模式的指南
| 场景 | 推荐模式 | 原因 |
| — | — | — |
| 前端实时更新 | updates | 带宽小,只传增量 |
| 完整状态监控 | values | 确保时刻掌握全局 |
| 聊天应用 Token 流 | messages | 原生支持,最优化 |
| 开发调试 | debug | 详尽信息,易定位问题 |
| 生产监控 | updates + debug | 组合使用 |
5.3 实战:构建实时聊天应用
import asyncio
from langgraph.graph import StateGraph, MessagesState, START, END
# 构建聊天图
builder = StateGraph(MessagesState)
def chatbot(state: MessagesState) -> dict:
response = llm.invoke(state["messages"])
return {"messages": [response]}
builder.add_node("chatbot", chatbot)
builder.add_edge(START, "chatbot")
builder.add_edge("chatbot", END)
graph = builder.compile()
# 异步流式调用(适合 Web 框架)
async def stream_chat(user_input: str):
"""流式返回聊天回复"""
async for chunk in graph.astream(
{"messages": [HumanMessage(user_input)]},
stream_mode="messages"
):
# 实时推送给前端
yield chunk
# 在 FastAPI 中使用
from fastapi.responses import StreamingResponse
@app.post("/chat")
async def chat_endpoint(msg: str):
return StreamingResponse(
stream_chat(msg),
media_type="application/x-ndjson"
)
# 前端接收
async function fetchChat(message) {
const response = await fetch("/chat?msg=" + message);
const reader = response.body.getReader();
while (true) {
const { done, value } = await reader.read();
if (done) break;
const text = new TextDecoder().decode(value);
// 实时显示 token
console.log(text);
}
}
第六部分:设计模式与最佳实践
6.1 工作流模式库
模式1:检查-处理-验证(Check-Process-Verify)
class CPVState(TypedDict):
user_input: str
validation_errors: List[str]
processed_data: dict
is_valid: bool
def check_input(state: CPVState) -> dict:
"""第一步:验证输入"""
errors = []
if not state["user_input"]:
errors.append("Input cannot be empty")
if len(state["user_input"]) > 1000:
errors.append("Input too long")
return {"validation_errors": errors, "is_valid": len(errors) == 0}
def process_data(state: CPVState) -> dict:
"""第二步:处理数据"""
if not state["is_valid"]:
return {"processed_data": {}}
data = parse_and_process(state["user_input"])
return {"processed_data": data}
def verify_output(state: CPVState) -> dict:
"""第三步:验证输出"""
if not state["processed_data"]:
return {"is_valid": False}
# 二次验证
return {"is_valid": True}
# 构建图
builder = StateGraph(CPVState)
builder.add_node("check", check_input)
builder.add_node("process", process_data)
builder.add_node("verify", verify_output)
builder.add_edge(START, "check")
def route_to_process(state: CPVState) -> str:
return "process" if state["is_valid"] else END
builder.add_conditional_edges("check", route_to_process)
builder.add_edge("process", "verify")
builder.add_edge("verify", END)
模式2:重试机制
class RetryState(TypedDict):
task: str
max_retries: int
retry_count: int
result: Optional[str]
error: Optional[str]
def execute_task(state: RetryState) -> dict:
"""执行任务,可能失败"""
try:
result = dangerous_operation(state["task"])
return {"result": result, "error": None}
except Exception as e:
return {"error": str(e), "result": None}
def decide_retry(state: RetryState) -> str:
"""决定是否重试"""
if state["error"] and state["retry_count"] < state["max_retries"]:
return "retry"
return END
# 构建图
builder = StateGraph(RetryState)
builder.add_node("execute", execute_task)
builder.add_edge(START, "execute")
builder.add_conditional_edges("execute", decide_retry)
def on_retry(state: RetryState) -> dict:
"""准备重试"""
return {"retry_count": state["retry_count"] + 1, "error": None}
builder.add_node("retry_prep", on_retry)
builder.add_edge("retry_prep", "execute")
# 路由:重试 -> retry_prep -> execute
builder.add_conditional_edges(
"execute",
lambda s: "retry" if s["error"] else END,
{"retry": "retry_prep"}
)
模式3:多代理协作(Supervisor Pattern)
class SupervisorState(TypedDict):
task: str
agent_results: Annotated[List[str], operator.add]
final_response: str
def supervisor(state: SupervisorState) -> str:
"""监督者决定选择哪个代理"""
if "data" in state["task"]:
return "data_agent"
elif "reasoning" in state["task"]:
return "reasoning_agent"
else:
return "general_agent"
def data_agent(state: SupervisorState) -> dict:
result = analyze_data(state["task"])
return {"agent_results": [f"DataAgent: {result}"]}
def reasoning_agent(state: SupervisorState) -> dict:
result = deep_reasoning(state["task"])
return {"agent_results": [f"ReasoningAgent: {result}"]}
def general_agent(state: SupervisorState) -> dict:
result = general_response(state["task"])
return {"agent_results": [f"GeneralAgent: {result}"]}
def synthesize(state: SupervisorState) -> dict:
"""综合所有代理的结果"""
combined = "\n".join(state["agent_results"])
final = synthesize_response(combined)
return {"final_response": final}
# 构建图
builder = StateGraph(SupervisorState)
builder.add_node("supervisor", supervisor)
builder.add_node("data_agent", data_agent)
builder.add_node("reasoning_agent", reasoning_agent)
builder.add_node("general_agent", general_agent)
builder.add_node("synthesize", synthesize)
builder.add_edge(START, "supervisor")
builder.add_conditional_edges(
"supervisor",
lambda x: x, # 直接返回节点名
{
"data_agent": "data_agent",
"reasoning_agent": "reasoning_agent",
"general_agent": "general_agent"
}
)
# 所有代理结果汇聚
builder.add_edge("data_agent", "synthesize")
builder.add_edge("reasoning_agent", "synthesize")
builder.add_edge("general_agent", "synthesize")
builder.add_edge("synthesize", END)
6.2 State 设计最佳实践
# ✅ DO:清晰分层的 State 设计
class RobustState(TypedDict):
# ===== 输入层 =====
user_input: str
user_id: str
session_id: str
# ===== 中间状态层 =====
parsed_intent: Optional[str] # 用户意图
extracted_entities: dict # 实体提取结果
current_step: str # 工作流步骤
# ===== 执行结果层 =====
messages: Annotated[List[BaseMessage], add_messages] # 对话历史
tool_results: Annotated[List[dict], operator.add] # 工具调用结果
# ===== 监控层 =====
processing_time: Annotated[float, lambda a, b: b] # 取最新
total_tokens: Annotated[int, lambda a, b: a + b] # 累加
errors: Annotated[List[str], operator.add] # 错误日志
# ===== 元数据层 =====
metadata: dict # 灵活存储额外信息
# ❌ DON'T:混乱的 State 设计
class MessyState(TypedDict):
data: str # 太模糊
x: int # 什么 x?
temp: dict # 临时什么?
stuff: list # 什么东西?
_private: str # 为什么用下划线?
CONSTANT: int # 为什么大写?
6.3 节点函数最佳实践
# ✅ DO:清晰、单一职责的节点
def analyze_sentiment(state: State) -> dict:
"""
单一职责:只做情感分析
文档完整、参数明确、返回值清晰
"""
try:
# 输入验证
if not state["messages"]:
return {"sentiment": "unknown", "score": 0}
# 核心逻辑
last_msg = state["messages"][-1].content
result = sentiment_model.predict(last_msg)
# 返回部分更新(关键!)
return {
"sentiment": result["label"],
"confidence": result["score"]
}
except Exception as e:
# 错误处理
return {
"sentiment": "error",
"errors": [str(e)]
}
# ❌ DON'T:职责不清、逻辑混乱的节点
def do_everything(state: State) -> dict:
"""
做太多事情,职责不清
"""
# 分析情感
sentiment = sentiment_model.predict(state["messages"][-1])
# 提取实体
entities = ner_model.extract(state["messages"][-1])
# 调用 API
api_result = requests.get("http://api.example.com")
# 保存数据库
db.insert(api_result)
# 发送通知
send_email(state["user_id"])
# 返回整个 state(太臃肿)
return state
6.4 错误处理与恢复最佳实践
from typing import Optional
class RobustState(TypedDict):
task: str
retries_left: int
last_error: Optional[str]
result: Optional[str]
def execute_with_fallback(state: RobustState) -> dict:
"""带回退策略的执行"""
try:
result = primary_method(state["task"])
return {"result": result, "last_error": None}
except PrimaryMethodError as e:
# 尝试备选方法
try:
result = fallback_method(state["task"])
return {"result": result, "last_error": None}
except FallbackMethodError as e2:
if state["retries_left"] > 0:
return {
"last_error": str(e2),
"retries_left": state["retries_left"] - 1
}
else:
return {
"result": default_response,
"last_error": f"Exhausted retries: {str(e2)}"
}
# 流程控制:根据错误决定下一步
def handle_result(state: RobustState) -> str:
if state["result"]:
return "success"
elif state["retries_left"] > 0:
return "retry"
else:
return "failure"
builder = StateGraph(RobustState)
builder.add_node("execute", execute_with_fallback)
builder.add_conditional_edges(
"execute",
handle_result,
{
"success": END,
"retry": "execute", # 重新执行
"failure": "handle_failure"
}
)
第七部分:生产级应用架构
7.1 可观测性与监控
import logging
from typing import Dict, Any
from datetime import datetime
class MonitoringMiddleware:
"""监控中间件:追踪节点执行时间、token 消耗等"""
def __init__(self):
self.metrics = {}
self.logger = logging.getLogger("langgraph")
def before_node(self, node_name: str, state: Dict[str, Any]):
"""节点执行前"""
self.metrics[node_name] = {
"start_time": datetime.now(),
"input_tokens": self._count_tokens(state)
}
self.logger.info(f"Node '{node_name}' started")
def after_node(self, node_name: str, state: Dict[str, Any]):
"""节点执行后"""
metrics = self.metrics[node_name]
metrics["end_time"] = datetime.now()
metrics["duration"] = (
metrics["end_time"] - metrics["start_time"]
).total_seconds()
metrics["output_tokens"] = self._count_tokens(state)
self.logger.info(
f"Node '{node_name}' completed in {metrics['duration']:.2f}s"
)
def _count_tokens(self, state: Dict[str, Any]) -> int:
"""估算 token 数"""
total = 0
for v in state.values():
if isinstance(v, str):
total += len(v) // 4
elif isinstance(v, list):
for item in v:
if hasattr(item, "content"):
total += len(item.content) // 4
return total
# 集成到图
monitor = MonitoringMiddleware()
# 包装节点函数
def monitored_node(original_node):
def wrapper(state: State) -> dict:
monitor.before_node(original_node.__name__, state)
result = original_node(state)
monitor.after_node(original_node.__name__, state)
return result
return wrapper
# 使用
@monitored_node
def process_data(state: State) -> dict:
return {"result": ...}
builder = StateGraph(State)
builder.add_node("process", process_data)
7.2 持久化与检查点
from langgraph.checkpoint.base import BaseCheckpointSaver
from langgraph.checkpoint.sqlite import SqliteSaver
import sqlite3
# 使用 SQLite 作为检查点存储
checkpoint = SqliteSaver.from_conn_string(
"file:langgraph_checkpoints.db"
)
graph = builder.compile(
checkpointer=checkpoint,
interrupt_before=["human_approval"], # 在这个节点前中断
auto_save=True # 自动保存检查点
)
# 恢复中断的执行
config = {
"configurable": {
"thread_id": "user_123" # 会话 ID
}
}
# 继续执行(从检查点恢复)
result = graph.invoke(
{"messages": [HumanMessage("approve")]},
config=config
)
# 查看执行历史
snapshot = graph.get_state(config)
print(f"Current values: {snapshot.values}")
print(f"Next nodes: {snapshot.next}")
7.3 上下文隔离与多租户
from contextlib import contextmanager
@dataclass
class TenantContext:
tenant_id: str
user_id: str
permissions: List[str]
rate_limit: int
class TenantContextMiddleware:
"""确保租户数据隔离"""
def validate_tenant_access(
self,
request: ModelRequest,
context: TenantContext
):
"""验证租户权限"""
if "admin_tools" in request.tools:
if "admin" not in context.permissions:
request.tools = [
t for t in request.tools
if t.name != "admin_tools"
]
# 检查速率限制
if not self._check_rate_limit(context):
raise RuntimeError(f"Rate limit exceeded for {context.tenant_id}")
def _check_rate_limit(self, context: TenantContext) -> bool:
# 实现速率限制逻辑
return True
# 使用
agent = create_agent(
model="gpt-4",
tools=all_tools,
middleware=[TenantContextMiddleware()],
context_schema=TenantContext
)
# 为每个租户创建隔离的执行上下文
def execute_for_tenant(tenant_id: str, user_id: str, query: str):
context = TenantContext(
tenant_id=tenant_id,
user_id=user_id,
permissions=get_user_permissions(user_id),
rate_limit=get_tenant_limit(tenant_id)
)
result = agent.invoke(
{"messages": [HumanMessage(query)]},
config={"configurable": {"context": context}}
)
return result
第八部分:常见陷阱与解决方案
8.1 状态爆炸
问题: 对话历史越来越长,导致 token 消耗不断增加。
# ❌ 不处理历史,无限增长
class BadState(TypedDict):
messages: Annotated[List[BaseMessage], add_messages]
# ✅ 定期裁剪历史
class SmartState(TypedDict):
messages: Annotated[List[BaseMessage], add_messages]
message_count: int
def trim_history(state: SmartState) -> dict:
"""每 20 条消息后进行裁剪"""
if len(state["messages"]) > 20:
# 保留最后 10 条
trimmed = state["messages"][-10:]
return {"messages": trimmed, "message_count": 10}
return {"message_count": len(state["messages"])}
# ✅ 使用摘要压缩
def summarize_history(state: SmartState) -> dict:
"""使用 LLM 生成摘要"""
if len(state["messages"]) > 30:
old_messages = state["messages"][:-5]
summary = llm.invoke(
f"Summarize this conversation:\n{old_messages}"
)
new_messages = [
SystemMessage(content=f"Summary: {summary.content}"),
*state["messages"][-5:]
]
return {"messages": new_messages}
return {}
8.2 条件边死循环
问题: 路由函数返回同一节点,导致无限循环。
# ❌ 危险:可能导致死循环
def bad_router(state: State) -> str:
if state["attempts"] < 3:
return "process" # 可能无限停留
return END
# ✅ 安全:显式限制循环次数
def safe_router(state: State) -> str:
if state["attempts"] < 3:
state["attempts"] += 1
return "process"
if state["result"] is None:
return "fallback"
return END
# ✅ 更安全:在编译时设置递归限制
graph = builder.compile(
config={"recursion_limit": 25} # 最多 25 个超步
)
8.3 Reducer 合并冲突
问题: 多个节点同时更新同一字段,顺序不可控。
# ❌ 顺序不可控
class BadState(TypedDict):
results: List[str] # 不用 Reducer
# 节点1、2、3 并行返回
# 节点1: {"results": ["A"]}
# 节点2: {"results": ["B"]}
# 节点3: {"results": ["C"]}
# 最终可能是 ["A", "B", "C"] 或 ["C", "B", "A"]
# ✅ 使用 Reducer 保证一致
class SmartState(TypedDict):
results: Annotated[List[str], operator.add]
# 或者用带序号的结构
result_map: dict # {"node_1": "A", "node_2": "B", "node_3": "C"}
def combine_results(state: SmartState) -> dict:
# 明确指定合并顺序
ordered = [
state["result_map"]["node_1"],
state["result_map"]["node_2"],
state["result_map"]["node_3"]
]
return {"final_results": ordered}
总结与关键要点
核心创新总结
| 特性 | v0.6 问题 | v1.0 方案 | 收益 |
| — | — | — | — |
| 上下文管理 | 配置地狱 | Context API | 代码简洁 60% |
| 工作流控制 | 黑盒循环 | StateGraph | 可视化、可调试 |
| 扩展机制 | 无 | 中间件 | 功能组合灵活 |
| 并行执行 | 不支持 | 超步并行 | 性能提升 3-5x |
| 类型安全 | 弱 | TypedDict/Pydantic | IDE 完整支持 |
实战最佳实践清单
- • ✅ 始终使用 TypedDict 定义 State
- • ✅ 每个节点只做一件事(单一职责)
- • ✅ 节点返回部分更新,不返回整个 State
- • ✅ 使用 Reducer 管理聚合字段(messages、results 等)
- • ✅ 为关键路由添加条件边,不是简单 IF-ELSE
- • ✅ 在编译时设置 recursion_limit 防止死循环
- • ✅ 使用 stream 而非 invoke 实现实时反馈
- • ✅ 中间件用于横切关注点(审计、安全、监控)
- • ✅ 生产环境必须配置 checkpointer 用于容错恢复
- • ✅ 定期裁剪或摘要历史消息,防止 token 爆炸
何时使用 LangGraph 1.0
推荐使用:
- • 复杂多步工作流(3+ 个步骤)
- • 需要用户交互或人工审批
- • 多 Agent 协作场景
- • 长对话需要记忆管理
- • 生产环境需要可观测性
可以用但不必要:
- • 简单的”输入 → LLM → 输出”
- • 单次对话,无状态
- • 原型快速验证
附录:完整工程示例
示例:企业级文档分析 Agent
from typing import TypedDict, Annotated, List, Optional
from dataclasses import dataclass
from langgraph.graph import StateGraph, START, END, Send
from langgraph.graph.message import add_messages
from langchain_openai import ChatOpenAI
from langchain_core.messages import BaseMessage, HumanMessage
import operator
# ===== State 定义 =====
class DocumentAnalysisState(TypedDict):
# 输入
document_id: str
document_content: str
analysis_type: str # "summary" | "sentiment" | "entities"
# 执行流程
messages: Annotated[List[BaseMessage], add_messages]
current_step: str
# 中间结果
analysis_results: Annotated[List[dict], operator.add]
# 监控
token_count: Annotated[int, lambda a, b: a + b]
errors: Annotated[List[str], operator.add]
# ===== 节点定义 =====
def validate_document(state: DocumentAnalysisState) -> dict:
"""验证文档"""
if not state["document_content"]:
return {"errors": ["Document is empty"]}
if len(state["document_content"]) > 100000:
return {"errors": ["Document too large"]}
return {"current_step": "validated"}
def analyze_document(state: DocumentAnalysisState) -> dict:
"""分析文档"""
llm = ChatOpenAI(model="gpt-4")
if state["analysis_type"] == "summary":
prompt = f"Summarize: {state['document_content'][:1000]}"
elif state["analysis_type"] == "sentiment":
prompt = f"Analyze sentiment: {state['document_content'][:1000]}"
else:
prompt = f"Extract entities: {state['document_content'][:1000]}"
response = llm.invoke(prompt)
return {
"analysis_results": [{
"type": state["analysis_type"],
"result": response.content
}],
"token_count": len(response.content) // 4
}
def format_output(state: DocumentAnalysisState) -> dict:
"""格式化输出"""
result = {
"document_id": state["document_id"],
"analysis": state["analysis_results"],
"token_count": state["token_count"]
}
return {"analysis_results": [result]}
# ===== 图构建 =====
builder = StateGraph(DocumentAnalysisState)
builder.add_node("validate", validate_document)
builder.add_node("analyze", analyze_document)
builder.add_node("format", format_output)
builder.add_edge(START, "validate")
def route_after_validation(state: DocumentAnalysisState) -> str:
if state["errors"]:
return END
return "analyze"
builder.add_conditional_edges("validate", route_after_validation)
builder.add_edge("analyze", "format")
builder.add_edge("format", END)
# 编译与执行
graph = builder.compile()
# 使用
result = graph.invoke({
"document_id": "doc_001",
"document_content": "This is a long document...",
"analysis_type": "summary"
})
print(f"Analysis result: {result['analysis_results']}")
print(f"Tokens used: {result['token_count']}")
这个完整指南涵盖了 LangGraph 1.0 的所有核心概念、API、设计模式和生产级实践,现在你可以开始构建复杂的 Agent 系统了!
免责声明:
本文所载程序、技术方法仅面向合法合规的安全研究与教学场景,旨在提升网络安全防护能力,具有明确的技术研究属性。
任何单位或个人未经授权,将本文内容用于攻击、破坏等非法用途的,由此引发的全部法律责任、民事赔偿及连带责任,均由行为人独立承担,本站不承担任何连带责任。
本站内容均为技术交流与知识分享目的发布,若存在版权侵权或其他异议,请通过邮件联系处理,具体联系方式可点击页面上方的联系我。
本文转载自:黄师傅的赛博dojo 黄师傅《LangGraph 1.0 完全指南:新概念、API、设计模式与最佳实践》