## 前言 如果你正在用 LangGraph 搭建 Agent,大概率遇到过这些场景: + Agent 陷入循环,疯狂调用工具把 Token 烧光 + 模型 API 偶尔抽风,整个链路直接崩掉 + 用户把身份证号、手机号一股脑塞进对话框 + 某个"删除生产数据"的操作,Agent 二话不说就执行了 这些问题,靠多写几行 `if-else` 或者往 System Prompt 里塞规则是解决不了的。你需要的是一个能**横切 Agent 执行流程**的治理层——这就是 Middleware。 本文把 LangGraph Agent Middleware 的完整知识体系梳理成一篇文章,从架构认知到底层原理,从内置中间件实战到自定义进阶。文中的所有示例围绕一个统一场景——**智能报表导出平台**——展开,模型选用国内开发者更熟悉的 DeepSeek 系列。 读完这篇,你应该能回答下面几个问题: + Middleware 在 Agent 执行循环的哪些位置介入? + 怎么给 Agent 装"刹车"(次数限制、人工确认、脱敏)? + 怎么让 Agent 长期稳定运行(重试降级、摘要压缩、上下文清理)? + 什么时候用内置 Middleware,什么时候自己写? + `request` / `handler` / `state` / `runtime` 这几个参数到底怎么用? --- ## 一、架构认知:Middleware 是什么,怎么运转 ### 1.1 为什么 Agent 需要被"管" 先看一个典型的 Agent 执行流程: B[🤖 调用大模型] B --> C{模型决定} C -->|调工具| D[⚙️ 执行工具] D --> E[📋 观察结果] E --> B C -->|直接回复| F[✅ 返回最终答案] --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/e9ca0629bff21f3991eedad528df50ae.svg) 这套循环本身没什么问题,但放到真实业务里,麻烦就来了。假设你在做一个智能报表导出 Agent,它有两个工具: + `check_export_permission`:查用户有没有导出权限 + `create_export_task`:创建导出任务,可能一次导出几十万行数据 Agent 能自己决策、自己调工具。但如果它决定反复查权限(浪费 Token)、在用户没说清楚场景时直接建导出任务(产生脏数据),或者用户消息里混着别人的手机号就送进了模型——这些都不是 Prompt 里多写两句话能拦住的事。 **Middleware 的定位就是:不改变 Agent 做业务的方式,但在它做业务的时候盯着、拦着、兜着底。** M1[before_agent] M1 --> M2[before_model] M2 --> M3[wrap_model_call] M3 --> LLM[大模型调用] LLM --> M4[after_model] M4 --> M5[wrap_tool_call] M5 --> Tool[工具执行] Tool --> M6[after_agent] M6 --> U end M1 -.->|"启动前检查"| M1 M2 -.->|"调用前日志"| M2 M3 -.->|"拦截/重试"| M3 M4 -.->|"响应后处理"| M4 M5 -.->|"工具管控"| M5 M6 -.->|"收尾清理"| M6 --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/953fa28923796d3c0ba1023a4d126fb4.svg) LangGraph 提供了**六大钩子点**,分别对应 Agent 执行流程的不同阶段: | 钩子 | 触发时机 | 典型用途 | | --- | --- | --- | | `before_agent` | Agent 启动前 | 初始化状态、检查前置条件 | | `before_model` | 每次调模型之前 | 审计日志、脱敏、状态检查 | | `wrap_model_call` | 包裹模型调用过程 | 重试、降级、模型切换 | | `after_model` | 模型返回之后 | 记录响应、检测 tool_calls、更新计数 | | `wrap_tool_call` | 包裹工具调用过程 | 权限校验、参数拦截、耗时统计 | | `after_agent` | Agent 执行结束后 | 清理状态、发送通知 | 一句话总结:**业务逻辑写在工具里,治理逻辑写在 Middleware 里。**职责分离是 Middleware 最核心的设计哲学。 ### 1.2 第一段 Middleware:从零到能观察到 Agent 在做什么 先搭一个最基础的 Agent——没有任何 Middleware,只有两个工具。场景是报表导出平台的权限校验: ```python from langchain.chat_models import init_chat_model from langchain.tools import tool # 使用 DeepSeek 模型,通过阿里云百炼平台接入 # init_chat_model 会自动适配 OpenAI 兼容接口 model = init_chat_model( model="deepseek-v4-flash", model_provider="openai", base_url="https://api.deepseek.com", api_key="your_api_key", ) # ---- 模拟数据 ---- # 实际项目中,这些数据来自数据库或 API EXPORT_RIGHTS = { "zhangsan": {"role": "finance", "region": "all"}, "lisi": {"role": "operation", "region": "east"}, } @tool def check_export_permission(username: str) -> dict: """查询用户是否有报表导出权限。 参数 username 为员工账号(英文名)。 """ user_info = EXPORT_RIGHTS.get(username) if not user_info: return {"username": username, "can_export": False, "reason": "用户不存在"} return { "username": username, "role": user_info["role"], "region": user_info["region"], "can_export": user_info["role"] in ("finance", "ops_manager"), } @tool def create_export_task( report_name: str, file_format: str, estimated_rows: int, reason: str ) -> dict: """创建报表导出任务。 report_name: 报表名称 file_format: 导出格式 (xlsx / csv) estimated_rows: 预计导出行数 reason: 导出原因 """ return { "task_id": "EXPORT-20260706-001", "report_name": report_name, "file_format": file_format, "estimated_rows": estimated_rows, "reason": reason, "status": "queued", } tools = [check_export_permission, create_export_task] ``` 这个 Agent 已经能干活了,但我们对它的内部行为一无所知。加上两个最轻量的 Middleware 来观察: ```python from langchain.agents import create_agent from langchain.agents.middleware import before_model, wrap_tool_call, AgentState @before_model def log_before_model(state: AgentState, runtime): """模型调用前执行:看一眼当前上下文里有多少条消息。""" # state["messages"] 是 Agent 当前累积的全部对话历史 # runtime 携带运行时环境信息(后面会详细讲) print(f"[审计] 准备调用模型,当前上下文消息数:{len(state['messages'])}") return None # 返回 None 表示不做任何修改 @wrap_tool_call def log_tool_call(request, handler): """工具调用前后各打一条日志,并记录耗时。""" tool_name = request.tool_call["name"] print(f"[审计] 开始执行工具:{tool_name}") # handler(request) 是"真正执行工具"的入口 # 不调用它,工具就不会执行 result = handler(request) print(f"[审计] 工具执行完毕:{tool_name}") return result # 组装 Agent,把 Middleware 列表传进去 agent = create_agent( model=model, tools=tools, middleware=[ log_before_model, # 排在前面:先记录状态 log_tool_call, # 排在后面:包裹工具调用 ], system_prompt=( "你是报表导出平台的智能助手。" "用户询问导出相关问题前,先调用 check_export_permission 确认权限。" "只有在用户明确提出导出需求时,才调用 create_export_task 创建任务。" ), ) # 跑一次看看 response = agent.invoke({ "messages": [ { "role": "user", "content": "我是 zhangsan,需要导出本月 east 区域的订单报表,大约 5000 行,xlsx 格式。", } ] }) print(response["messages"][-1].content) ``` 运行时你会看到类似这样的输出: ```plain [审计] 准备调用模型,当前上下文消息数:1 [审计] 开始执行工具:check_export_permission [审计] 工具执行完毕:check_export_permission [审计] 准备调用模型,当前上下文消息数:3 [审计] 开始执行工具:create_export_task [审计] 工具执行完毕:create_export_task [审计] 准备调用模型,当前上下文消息数:5 ``` Agent 的业务能力没变,但现在你能清楚地看到:模型被调了几次、每次调了哪些工具、消息在什么时候膨胀。这就是 Middleware 的第一个价值——**可观测性**,不侵入业务代码。 ### 1.3 深入 `request` 与 `handler` 写 `@wrap_tool_call` 或 `@wrap_model_call` 时,你一定会遇到这两个参数。它们的直觉含义是"拦住一次调用,看看参数,再决定要不要继续"。但很多人卡在对这两个对象的具体结构不熟悉上。 #### 1.3.1 工具级:`ToolCallRequest` 当你在 `@wrap_tool_call` 里拿到 `request` 时,它是一个 `ToolCallRequest` 对象,包含以下关键字段: | 字段 | 含义 | 用法举例 | | --- | --- | --- | | `request.tool_call["name"]` | 模型决定调哪个工具 | `if tool_name == "create_export_task"` | | `request.tool_call["args"]` | 模型生成的调用参数 | 校验 `estimated_rows` 是否超标 | | `request.tool_call["id"]` | 本次调用的唯一 ID | 日志追踪 | | `request.tool` | 工具对象本身 | 读取工具的 docstring | | `request.state` | Agent 当前累积状态 | 查看已有消息数 | | `request.runtime` | 运行时上下文 | 读取从业务层传入的 `context` | `handler`** 是一个可调用对象**——`handler(request)` 执行之后,真正的工具才会运行。这意味着你可以在调用 `handler` 之前做校验、修改参数,甚至直接跳过它。 下面是一个实用的例子:在工具执行前拦截超标请求,不调用真正的导出工具,而是返回伪造的 `ToolMessage` 让 Agent 知道被拦了: ```python from langgraph_core.messages import ToolMessage @wrap_tool_call def guard_export_scale(request, handler): """超过 10 万行的导出请求直接拦截,不进入执行队列。""" tool_name = request.tool_call["name"] args = request.tool_call.get("args", {}) # 只拦截 create_export_task if tool_name == "create_export_task" and args.get("estimated_rows", 0) > 100_000: # 不调 handler,直接返回 ToolMessage 给模型 # 模型收到这条消息后,会告知用户"太大了,换个小范围" return ToolMessage( content=f"导出被拦截:预计行数 {args.get('estimated_rows')} 超过上限 100000,请缩小筛选范围后重试。", tool_call_id=request.tool_call["id"], ) # 正常放行 return handler(request) ``` 四种 `handler` 用法模式: Q{怎么处理 handler?} Q -->|"模式1: 透明放行"| A["return handler(request)"] Q -->|"模式2: 拦截替换"| B["return ToolMessage(...)"] Q -->|"模式3: 重试覆盖"| C["handler(request.override(model=backup))"] Q -->|"模式4: 链式修改"| D["修改 request 后调 handler"] --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/62f94cf516ce1e70acfd1c64366f453a.svg) > ⚠️ **重要**:不调 `handler` = 原始操作从未发生;调两次 `handler` = 操作执行两次——对于发邮件、扣款这类操作,这可能造成严重后果。 > #### 1.3.2 模型级:`ModelRequest` `@wrap_model_call` 里的 `request` 是 `ModelRequest`,结构更"厚": | 字段 | 含义 | | --- | --- | | `request.messages` | 即将发送给模型的消息列表——就是上文说的"上下文" | | `request.model` | 当前使用的模型对象 | | `request.tools` | 模型可用的工具列表 | | `request.model_settings` | 模型参数(temperature 等) | | `request.state` / `request.runtime` | 同工具级 | 这里有一个很关键的用法——`request.override(model=...)`,可以在不重建 Agent 的情况下临时切换模型: ```python @wrap_model_call def local_first_then_cloud(request, handler): """先用本地模型,失败 3 次后切云端模型。""" # 先用本地 Ollama 模型(省成本) try: return handler(request) # 使用 Agent 默认模型 except Exception: # 本地挂了,切到 DeepSeek 云端 print("[降级] 本地模型不可用,切换到云端 DeepSeek") return handler(request.override(model=cloud_model)) ``` ```python # 一个典型的 Agent 运行流程: 用户输入 → [wrap_model_call 拦截] → LLM 思考 → 决定调用工具 get_weather → [wrap_tool_call 拦截] → 执行 get_weather 工具 → [wrap_model_call 拦截] → LLM 再次思考 → 生成最终回复 → 返回给用户 ``` + `**wrap_model_call**` 在模型每次被调用时触发,包括首次决策和后续的反思轮次。 + `**wrap_tool_call**` 只在模型明确要求调用工具时触发。 ### 1.4 深入 `state` 与 `runtime` 如果说 `request` / `handler` 关注的是"**这一次调用**",那 `state` / `runtime` 关注的就是"**Agent 当前整体处于什么状态**"。 H[handler] end subgraph 全局运行视角 S[state: 累积消息/计数] --- RT[runtime: 外部上下文] end --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/13c6d865e0270f532a4306f1f95c86e9.svg) #### 1.4.1 `state`:Agent 的内部状态 `state`(类型 `AgentState`)是 Agent 在运行过程中不断累积的数据。最核心的字段是 `state["messages"]`——所有对话历史和工具结果都在里面。 你可以扩展 `AgentState`,添加自定义字段来追踪业务指标: ```python from typing import TypedDict from typing_extensions import NotRequired from langchain.agents.middleware import AgentState, before_model, after_model class ExportAgentState(AgentState): """扩展默认状态,增加模型调用计数和大任务标记。""" # AgentState 已经自带 messages 字段,这里只加新字段 # NotRequired 表示可以不传,Middleware 内部自己维护 model_call_count: NotRequired[int] blocked_requests: NotRequired[int] # 被拦截的请求次数 @before_model(state_schema=ExportAgentState) def audit_before_model(state: ExportAgentState, runtime): """每次调模型前,看一眼当前统计。""" count = state.get("model_call_count", 0) blocked = state.get("blocked_requests", 0) print(f"[统计] 第 {count + 1} 次调模型 | 已拦截 {blocked} 次") return None @after_model(state_schema=ExportAgentState) def update_stats(state: ExportAgentState, runtime): """模型返回后,把计数器 +1。""" # after_model 返回 dict 可以直接更新 AgentState return {"model_call_count": state.get("model_call_count", 0) + 1} ``` 关键点:`@before_model` / `@after_model` 需要指定 `state_schema`,告诉 LangGraph 你的自定义状态有哪些字段。`@after_model` 返回的 `dict` 会被浅合并进 `state`。 #### 1.4.2 `runtime`:外部注入的上下文 `runtime.context` 存放的是**不属于对话本身、但从外部传入的运行时信息**——比如请求 ID、用户角色、业务标识。 ```python from langgraph.runtime import Runtime class RunContext(TypedDict): """定义 runtime.context 的结构。""" request_id: str # 用于日志追踪 user_role: str # 当前用户的角色 tenant_id: str # 租户标识(多租户场景) @before_model(state_schema=ExportAgentState) def inject_context(state: ExportAgentState, runtime: Runtime[RunContext]): """从 runtime.context 读取业务信息,用于日志关联。""" ctx = runtime.context or {} print( f"[上下文] request_id={ctx.get('request_id')} " f"user_role={ctx.get('user_role')} " f"tenant={ctx.get('tenant_id')}" ) return None # 创建 Agent 时声明 context 结构 agent = create_agent( model=model, tools=tools, middleware=[audit_before_model, inject_context, update_stats], state_schema=ExportAgentState, # 声明自定义状态 context_schema=RunContext, # 声明上下文结构 system_prompt="你是报表导出平台的智能助手。", ) # 调用时传入 context result = agent.invoke( {"messages": [{"role": "user", "content": "导出本月的订单报表"}]}, context={ "request_id": "req-20260706-001", "user_role": "finance", "tenant_id": "t-1234", }, ) ``` > **重要**:`runtime.context` **不会自动让模型看到**。如果你想让它影响模型行为,需要在 Middleware 里显式地把信息写入 state 或拼入 prompt——这正是 `@dynamic_prompt` 的用武之地(见第四章)。 > --- ## 二、安全风控:给 Agent 装上刹车和护栏 线上 Agent 有三类典型风险:**调用失控**(无限循环烧 Token)、**操作越权**(关键动作未经确认)、**隐私泄露**(PII 直接喂给模型)。这一章用三个内置 Middleware 分别解决。 ### 2.1 调用次数限制:防住"无穷循环" 假设 Agent 要查一个异步导出任务的进度。任务可能跑好几分钟,Agent 会不停地轮询。但如果没有限制,它能在 30 秒内调用 200 次模型——Token 账单直接爆炸。 `ModelCallLimitMiddleware` 和 `ToolCallLimitMiddleware` 就是专门干这个的。 #### ModelCallLimitMiddleware:限制模型调用次数 ```python from langchain.agents.middleware import ModelCallLimitMiddleware # 模拟一个"前 4 次返回处理中,第 5 次返回完成"的进度查询工具 export_progress = {"task_001": 0} @tool def check_export_progress(task_id: str) -> str: """查询导出任务的进度。""" export_progress[task_id] = export_progress.get(task_id, 0) + 1 attempt = export_progress[task_id] if attempt < 5: return f"第 {attempt} 次查询:任务仍在处理中,请稍候。" return f"第 {attempt} 次查询:导出完成,文件已生成。" # 不加限制:Agent 会一直查到第 5 次 # 加上 ModelCallLimitMiddleware(run_limit=3):最多调 3 次模型就强制结束 agent = create_agent( model=model, tools=[check_export_progress], middleware=[ ModelCallLimitMiddleware( run_limit=3, # 单次 invoke 最多调 3 次模型 exit_behavior="end", # 到达上限后尝试优雅结束(生成总结) ), ], system_prompt=( "你是导出任务进度查询助手。" "当用户查询任务进度时,调用 check_export_progress。" "如果任务还在处理中,继续查询直到完成。" ), ) result = agent.invoke({ "messages": [{"role": "user", "content": "帮我查 task_001 的导出进度,持续查到完成为止。"}] }) # 输出类似:"Model call limits exceeded: run limit (3/3)" # Agent 被强制刹车,不会无限循环 print(result) ``` #### ToolCallLimitMiddleware:限制特定工具的调用次数 某些场景下你不想限制整个 Agent 的模型调用次数,只想限制**某一个工具**——比如进度查询工具被频繁轮询。`ToolCallLimitMiddleware` 只针对指定工具生效: ```python from langchain.agents.middleware import ToolCallLimitMiddleware agent = create_agent( model=model, tools=[check_export_progress], middleware=[ ToolCallLimitMiddleware( tool_name="check_export_progress", # 只限制这个工具 run_limit=2, # 最多调 2 次 exit_behavior="continue", # 达到上限后通知模型继续(不抛异常) ), ], system_prompt=( "你是导出任务进度查询助手。" "当 check_export_progress 返回限制信息时,告知用户当前进度并建议稍后再查。" ), ) # exit_behavior="continue" 的效果: # 工具达到上限后,Middleware 返回一条 ToolMessage 告知模型"已达调用上限", # 模型可以据此生成友好的回复,而不是直接报错 ``` > **两个 Middleware 可以混用**:同时配置 `ModelCallLimitMiddleware`(全局上限)和 `ToolCallLimitMiddleware`(单工具上限),Agent 会在任一限制触发时终止或被拦截。 > #### run_limit vs thread_limit | 参数 | 作用域 | 适用场景 | | --- | --- | --- | | `run_limit` | 单次 `invoke()` | 防止一次请求中过度调用 | | `thread_limit` | 整个对话线程(需 checkpointer) | 跨多轮对话的累计限制 | `thread_limit` 需要配合 `checkpointer` 使用,因为跨轮追踪需要持久化状态: ```python from langgraph.checkpoint.memory import InMemorySaver agent = create_agent( model=model, tools=[check_export_progress], middleware=[ ModelCallLimitMiddleware( run_limit=10, thread_limit=50, # 整个对话生命周期最多 50 次模型调用 exit_behavior="end", ), ], checkpointer=InMemorySaver(), # 必需:跨轮追踪需要存储 system_prompt="你是导出任务进度查询助手。", ) ``` #### `exit_behavior` 的三种模式 | 值 | 行为 | 适用场景 | | --- | --- | --- | | `"end"` | 优雅结束,生成兜底回复 | 大多数场景 | | `"error"` | 抛出异常,上层捕获处理 | 需要明确感知限流事件 | | `"continue"` | (仅工具限制)返回限制信息给模型,让模型决定怎么说 | 模型可以继续回复用户的场景 | ### 2.2 关键操作人工确认:Human-in-the-Loop 有些操作 Agent 不应该自己决定——比如创建一个可能消耗大量资源的导出任务。`HumanInTheLoopMiddleware` 让 Agent 在执行特定工具前**暂停并等待人类审批**。 >A: ⏸️ 暂停!需要人工审批 M->>U: 请确认:创建导出任务? U->>M: ✅ approve(或 ✏️ edit / ❌ reject) M->>T: 放行(或修改参数后放行) T->>A: 返回执行结果 A->>U: 任务已创建 --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/f8c955f1417539fabab1b1cb0d73d5c3.svg) #### 实战:为危险操作加上确认环节 ```python from typing import Literal from pydantic import BaseModel, Field from langchain.agents.middleware import HumanInTheLoopMiddleware from langgraph.checkpoint.memory import InMemorySaver # 第一步:用 Pydantic 定义工具的输入结构 # 结构化的参数让模型生成更准确,也方便人工审批时查看 class ExportTaskInput(BaseModel): """导出任务的参数结构。模型会按照这个 schema 生成参数。""" report_name: str = Field(description="报表名称,如 'east_region_orders_202607'") file_format: Literal["xlsx", "csv"] = Field(description="导出格式") estimated_rows: int = Field(description="预计导出行数") reason: str = Field(description="导出原因,用于审计") priority: Literal["low", "normal", "high"] = Field( default="normal", description="优先级" ) # 第二步:安全工具不加 args_schema(自动执行) @tool def check_export_permission(username: str) -> dict: """查询用户导出权限。安全操作,无需审批。""" return {"username": username, "can_export": True, "region": "all"} # 第三步:危险工具加上 args_schema,配合 HITL 拦截 @tool(args_schema=ExportTaskInput) def create_export_task( report_name: str, file_format: str, estimated_rows: int, reason: str, priority: str = "normal", ) -> dict: """创建报表导出任务。这是一个高风险操作,需要人工审批。""" print( f"[导出] 创建任务:{report_name} | {file_format} | " f"{estimated_rows} 行 | 优先级 {priority}" ) return { "task_id": "EXPORT-20260706-002", "report_name": report_name, "file_format": file_format, "estimated_rows": estimated_rows, "status": "queued", } # 第四步:创建带 HITL 的 Agent checkpointer = InMemorySaver() # HITL 必需:暂停后需要从这里恢复状态 agent = create_agent( model=model, tools=[check_export_permission, create_export_task], middleware=[ HumanInTheLoopMiddleware( interrupt_on={ # 安全工具:不中断,自动执行 "check_export_permission": False, # 危险工具:中断,提供三种审批选项 "create_export_task": { "allowed_decisions": ["approve", "edit", "reject"], }, }, ), ], checkpointer=checkpointer, system_prompt=( "你是报表导出平台的智能助手。" "用户查询权限时,直接调用 check_export_permission。" "只在用户明确要求创建导出任务时调用 create_export_task。" ), ) # 第五步:第一次 invoke——会被 HITL 拦截 config = {"configurable": {"thread_id": "export-001"}} result = agent.invoke( { "messages": [ { "role": "user", "content": "我是 zhangsan,导出 east 区本月订单报表,xlsx,约 5000 行,月度对账用。", } ] }, config=config, ) # result 里会包含中断信息,UI 层可以据此展示审批界面 print("Agent 已暂停,等待审批...") ``` #### 三种审批操作 ```python # 审批通过:工具以原始参数执行 agent.invoke( Command(resume={"decisions": [{"type": "approve"}]}), config=config, ) # 编辑后通过:修改参数再执行(比如把 estimated_rows 从 50000 改成 5000) agent.invoke( Command(resume={ "decisions": [ { "type": "edit", "edited_action": { "name": "create_export_task", "args": { "report_name": "east_region_orders_202607", "file_format": "xlsx", "estimated_rows": 5000, # 人工修正 "reason": "月度对账", "priority": "high", }, }, } ] }), config=config, ) # 拒绝:工具不执行,拒绝原因会反馈给模型 agent.invoke( Command(resume={ "decisions": [ { "type": "reject", "message": "审批被驳回:本月对账已由系统自动完成,无需手动导出。", } ] }), config=config, ) ``` > ⚠️ **HITL 的三个要点**: > > 1. **必须有 checkpointer**——Agent 暂停后需要持久化状态,`InMemorySaver` 适合开发调试,生产环境建议用 `SqliteSaver` 或 `PostgresSaver` > 2. `thread_id`** 是恢复的钥匙**——创建任务和审批操作必须用同一个 `thread_id` > 3. **不是所有工具都要审批**——只对真正高危的操作加 HITL,给每个操作都加确认会让用户体验很差 > ### 2.3 敏感信息自动脱敏 用户经常不经意地在消息里放敏感信息:"帮我查一下导出记录,我的手机号 13800138000,身份证 110101199003070013"。 这些信息不应该原样送进模型。`PIIMiddleware` 在消息进入模型**之前**做检测和替换: P[🔍 PIIMiddleware] P -->|检测通过| M[🤖 模型] P -->|检测到 PII| R{策略?} R -->|redact| R1["替换为 [REDACTED_PHONE]"] R -->|hash| R2["替换为一致的 hash 值"] R -->|mask| R3["部分遮盖如 138****8000"] R -->|block| R4["直接拒绝,抛出异常"] R1 --> M R2 --> M R3 --> M --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/ca4a93d5a95984fb5cc6f78dac7e5f68.svg) #### 内置检测类型 + 自定义正则 LangGraph 内置支持 email、信用卡号、IP 地址、MAC 地址、URL 等。中文场景常用的手机号、身份证号需要用自定义正则: ```python from langchain.agents.middleware import PIIMiddleware # 中国大陆手机号:1 开头,第二位 3-9,共 11 位 PHONE_PATTERN = r"(? `block` 抛出异常后会中断整个 Agent 执行流。如果你的应用需要优雅降级(比如告诉用户"请勿发送敏感信息"),需要在 `agent.invoke()` 外层 `try/except PIIDetectionError`。 > #### 不止输入:工具返回和输出也要查 ```python agent = create_agent( model=model, tools=tools, middleware=[ PIIMiddleware( "email", strategy="redact", apply_to_input=True, # 用户输入先过滤 apply_to_tool_results=True, # 工具返回值也过滤 apply_to_output=True, # 模型最终输出再检查一遍 ), ], system_prompt="你是报表导出客服助手。", ) ``` > ⚠️ **注意**:`PIIMiddleware` 只保护 **Agent 消息管道**。你自己工具里打的日志、调用的外部 API,不在它的保护范围内——那些需要另外做脱敏。 > --- ## 三、稳定性与上下文治理:让 Agent 长期可靠运行 Agent 在生产环境跑久了,会遇到两类问题:**外部依赖偶尔失败**(模型超时、接口抖动)和**内部状态持续膨胀**(对话越来越长、旧工具结果堆积)。这一章用三组 Middleware 解决。 ### 3.1 失败重试与模型降级 网络抖动、API 限流、模型服务临时不可用——这些在生产环境是常态。`ToolRetryMiddleware` 和 `ModelRetryMiddleware` 处理"重试"问题,`ModelFallbackMiddleware` 处理"降级"问题。 B{成功?} B -->|是| C[继续执行] B -->|否| D{还有重试次数?} D -->|是| E[等待延迟后重试] E --> B D -->|否| F{有备用模型?} F -->|是| G[切换备用模型] F -->|否| H[按 on_failure 策略处理] --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/0b6975d8057a98e1a6503223a113a080.svg) #### 工具重试:处理间歇性超时 ```python from langgraph.agents.middleware import ToolRetryMiddleware # 模拟一个前两次超时、第三次成功的导出接口 QUERY_ATTEMPTS = {} @tool def fetch_export_history(user_id: str) -> dict: """查询用户的历史导出记录。""" QUERY_ATTEMPTS[user_id] = QUERY_ATTEMPTS.get(user_id, 0) + 1 attempt = QUERY_ATTEMPTS[user_id] if attempt < 3: # 前两次模拟网络超时 raise TimeoutError(f"查询接口超时(第 {attempt} 次)") # 第三次成功 return { "user_id": user_id, "records": [ {"task_id": "EXP-001", "report": "月度订单", "time": "2026-07-01"}, {"task_id": "EXP-002", "report": "客户分析", "time": "2026-07-05"}, ], } agent = create_agent( model=model, tools=[fetch_export_history], middleware=[ ToolRetryMiddleware( tools=["fetch_export_history"], # 只对这个工具启用重试 max_retries=2, # 额外重试 2 次(共 3 次机会) retry_on=(TimeoutError,), # 只在超时时重试,权限错误不重试 initial_delay=0.2, # 第一次重试前等 0.2 秒 max_delay=1.0, # 重试间隔上限 1 秒(指数退避) on_failure="continue", # 重试耗尽后让模型继续(不抛异常) ), ], system_prompt="你是报表导出平台的查询助手。根据查询结果回答用户。", ) ``` 关键设计理念:**重试逻辑不写进工具函数,由 Middleware 统一处理。** 你的工具只管"查数据",超时了怎么办、重试多少次、间隔多久——这些是 Middleware 该操心的事。 #### 模型重试与降级:本地模型挂了切云端 ```python from langchain.agents.middleware import ModelRetryMiddleware, ModelFallbackMiddleware # 本地模型:Ollama + Qwen3,成本低但在高负载下可能不可用 local_model = init_chat_model("ollama:qwen3:14b") # 云端模型:DeepSeek,稳定但要花钱 cloud_model = init_chat_model( model="deepseek-v4-flash", model_provider="openai", base_url="https://api.deepseek.com", api_key="your_api_key_here", ) agent = create_agent( model=local_model, # 默认用本地模型 tools=tools, middleware=[ # 第一层:同一个模型失败后重试 ModelRetryMiddleware( max_retries=3, # 最多重试 3 次 retry_on=(Exception,), # 任何异常都重试 on_failure="continue", ), # 第二层:重试也失败了,切备用模型 ModelFallbackMiddleware( cloud_model, # 备用的云端模型 ), ], system_prompt="你是报表导出平台的智能助手。", ) ``` > **选型建议**:重试适合处理临时性抖动(网络波动、短暂限流);降级适合处理持续性不可用(服务宕机、额度耗尽)。不要把验证错误(如参数校验失败)放进重试范围——重试一万次也不会通过。 > ### 3.2 长对话自动摘要:别让上下文吃掉你的 Token 多轮对话中,用户可能前前后后提了十几次需求。所有历史消息都塞进上下文,模型窗口再大也扛不住——而且大部分旧消息对当前回答帮助不大。 `SummarizationMiddleware` 的思路很朴素:**超过阈值后,把旧消息压缩成一段摘要,只保留最近几条原文。** B{达到触发阈值?} B -->|否| A B -->|是| C[取旧消息 → 调摘要模型压缩] C --> D[摘要 + 最近 N 条原文 → 新上下文] --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/b9c4231915c73c0d75d22bda433b977f.svg) ```python from langchain.agents.middleware import SummarizationMiddleware agent = create_agent( model=cloud_model, tools=[], middleware=[ SummarizationMiddleware( model=model, # trigger=("messages", 8), # 消息数达到 8 条时触发 keep=("messages", 4), # 保留最近 4 条原文 # 自定义摘要提示词:强调保留"不做的事情" summary_prompt=( "请将以下对话历史压缩为一份结构化摘要。必须包含:\n" "1. 用户的最终目标\n" "2. 已确认的需求要点\n" "3. 用户明确表示不需要的功能(否定约束)\n" "4. 当前待解决的问题\n" "摘要不超过 300 字。" ), ), ], system_prompt="你是报表导出平台的需求分析助手。", ) ``` **关键参数说明**: | 参数 | 可选值 | 说明 | | --- | --- | --- | | `trigger` | `("messages", N)` / `("tokens", N)` / `("fraction", 0.8)` | 触发摘要的条件 | | `keep` | `("messages", N)` / `("tokens", N)` / `("fraction", 0.25)` | 保留多少最近内容不压缩 | | `model` | 任何 ChatModel | 做摘要的模型,可以和主模型不同 | | `summary_prompt` | 自定义 prompt 字符串 | 默认摘要有时会漏掉"不做的事" | > **常见坑**:默认摘要往往会**忽略否定信息**。比如用户说过"不支持 PDF 导出"、"移动端不做",摘要里很可能就丢了。自定义 `summary_prompt` 时,显式要求保留这类约束。 > ### 3.3 清理旧工具结果:别让过期的数据干扰模型 `SummarizationMiddleware` 压缩的是**聊天历史**,但还有一种膨胀——**工具返回值**。 Agent 一个任务调 4-5 次工具很正常。每次工具返回几百上千字,几轮下来上下文里塞满了旧结果。但模型做下一轮决策时,通常只需要最近一次工具的结果。`ContextEditingMiddleware` + `ClearToolUsesEdit` 解决的就是这个问题。 ```python from langchain.agents.middleware import ContextEditingMiddleware, ClearToolUsesEdit agent = create_agent( model=model, tools=tools, middleware=[ ContextEditingMiddleware( edits=[ ClearToolUsesEdit( trigger=120, # token 数超过 120 时触发 keep=1, # 保留最近 1 次工具结果 placeholder="[旧工具结果已被清理,节省上下文空间]", exclude_tools=("critical_audit_log",), # 审计日志不清理 clear_at_least=50, # 每次至少清理 50 token clear_tool_inputs=True, # 同时截断 AIMessage 里的 tool_call args ), ], ), ], system_prompt="你是报表导出平台的智能助手。", ) ``` **关键补充参数**: | 参数 | 含义 | 说明 | | --- | --- | --- | | `clear_tool_inputs` | 是否同时清理 AIMessage 中的工具调用参数 | 设 `True` 后,旧 AIMessage 里的 `tool_calls[].args` 也会被截断,进一步节省空间 | | `token_count_method` | Token 计数策略 | `"approximate"`(默认):用字符数估算,零开销;`"model"`:调模型精确计数,更准确但有额外成本 | | `clear_at_least` | 触发后最少清理多少 token | 防止触发清理但实际回收空间太小(比如刚好跨过阈值 1 token) | #### 摘要 vs 清理:什么时候用哪个 | 维度 | SummarizationMiddleware | ContextEditingMiddleware | | --- | --- | --- | | 处理对象 | 对话消息(Human + AI) | 工具调用结果(ToolMessage) | | 处理方式 | 压缩成摘要文本 | 替换为占位符 | | 适合场景 | 历史对话信息密度低 | 工具返回数据量大且时效性强 | | 组合使用 | 可以,先摘要再清理,各自负责不同维度 | | --- ## 四、自定义进阶:写出自己的 Middleware 内置 Middleware 覆盖了大部分通用场景,但每个项目的业务规则千差万别。比如:"导出任务超过 10 万行时自动降级为异步"、"模型返回后检查是否包含敏感字段名"——这些就需要自己动手了。 LangGraph 提供了**两条路径**:装饰器模式和类模式。 ### 4.1 装饰器模式:轻量级,适合单个 Hook 四个装饰器分别对应四个钩子点: | 装饰器 | 适用场景 | | --- | --- | | `@before_model` | 调用前作日志、状态检查、上下文注入 | | `@after_model` | 调用后记录结果、更新计数、触发后续动作 | | `@wrap_model_call` | 包裹模型调用,实现重试、降级、超时 | | `@wrap_tool_call` | 包裹工具调用,实现参数校验、权限拦截、耗时统计 | 下面是一个完整的自定义 Middleware 组,覆盖报表导出平台的业务规则: ```python import time from langchain.agents.middleware import ( before_model, after_model, wrap_model_call, wrap_tool_call, AgentState, ) from langchain_core.messages import ToolMessage # ---- Middleware 1:模型调用前记录审计信息 ---- @before_model def audit_before_model(state: AgentState, runtime): """在模型调用前,记录当前上下文规模和来源用户。""" ctx = runtime.context or {} print( f"[审计] 模型调用 #{state.get('call_count', 0) + 1} | " f"消息数 {len(state['messages'])} | " f"用户 {ctx.get('username', 'unknown')}" ) return None # ---- Middleware 2:模型返回后更新调用计数 ---- @after_model def update_call_counter(state: AgentState, runtime): """模型返回后,调用计数 +1。""" return {"call_count": state.get("call_count", 0) + 1} # ---- Middleware 3:包裹模型调用,实现本地重试 + 云端降级 ---- local_failures = 0 MAX_LOCAL_FAILURES = 3 @wrap_model_call def retry_local_then_cloud(request, handler): """先用本地 Ollama 模型,失败 N 次后切到云端 DeepSeek。""" global local_failures try: result = handler(request) # 尝试用本地模型 local_failures = 0 # 成功后重置计数器 return result except Exception as e: local_failures += 1 if local_failures < MAX_LOCAL_FAILURES: print(f"[降级] 本地模型失败({local_failures}/{MAX_LOCAL_FAILURES}),重试中...") raise # 重新抛出,让 ModelRetryMiddleware 处理 print(f"[降级] 本地模型不可用,切换到云端 DeepSeek") # request.override 临时换模型,不影响 Agent 默认配置 return handler(request.override(model=cloud_model)) # ---- Middleware 4:包裹工具调用,实现业务规则拦截 ---- @wrap_tool_call def guard_export_tool(request, handler): """在 create_export_task 执行前,做业务规则校验。""" tool_name = request.tool_call["name"] if tool_name == "create_export_task": args = request.tool_call.get("args", {}) estimated_rows = args.get("estimated_rows", 0) # 规则 1:超过 100 万行直接拒绝 if estimated_rows > 1_000_000: return ToolMessage( content=( f"导出被拦截:预计行数 {estimated_rows} 超过系统上限 1000000。" "建议:缩小时间范围或分批次导出。" ), tool_call_id=request.tool_call["id"], ) # 规则 2:超过 10 万行发出警告但放行(打印日志供运维关注) if estimated_rows > 100_000: print( f"[警告] 大导出任务:{args.get('report_name')} | " f"{estimated_rows} 行 | 用户可能需要等待较长时间" ) # 规则 3:记录工具调用耗时 start = time.time() result = handler(request) elapsed = time.time() - start print(f"[性能] {tool_name} 执行耗时 {elapsed:.2f}s") return result # 非目标工具直接放行 return handler(request) ``` ### 4.2 类模式:可配置、可复用的 Middleware 装饰器适合快速写一个 hook,但如果你的 Middleware 需要配置参数(比如"不同角色的用户有不同的导出上限"),类模式更合适——继承 `AgentMiddleware`。 ```python from langchain.agents.middleware import AgentMiddleware class ExportGovernanceMiddleware(AgentMiddleware): """报表导出治理中间件:按角色限制导出规模。 用法: ExportGovernanceMiddleware( role_limits={"finance": 500_000, "operation": 100_000, "default": 50_000} ) """ def __init__(self, role_limits: dict = None): super().__init__() self.role_limits = role_limits or {"default": 100_000} def before_model(self, state, runtime): """每次模型调用前,注入当前用户的权限信息到日志。""" ctx = runtime.context or {} role = ctx.get("user_role", "default") limit = self.role_limits.get(role, self.role_limits["default"]) print(f"[治理] 用户角色 {role},导出上限 {limit} 行") return None def wrap_tool_call(self, request, handler): """在创建导出任务前,检查是否超过该角色的导出上限。""" if request.tool_call["name"] != "create_export_task": return handler(request) ctx = request.runtime.context or {} role = ctx.get("user_role", "default") limit = self.role_limits.get(role, self.role_limits["default"]) args = request.tool_call.get("args", {}) estimated_rows = args.get("estimated_rows", 0) if estimated_rows > limit: return ToolMessage( content=( f"导出被拦截:预计 {estimated_rows} 行," f"超过 {role} 角色的上限 {limit} 行。" "请联系上级审批或缩小导出范围。" ), tool_call_id=request.tool_call["id"], ) return handler(request) # 使用 agent = create_agent( model=model, tools=tools, middleware=[ ExportGovernanceMiddleware( role_limits={ "finance": 500_000, "operation": 100_000, "default": 50_000, } ), ], context_schema=RunContext, system_prompt="你是报表导出平台的智能助手。", ) ``` > **装饰器 vs 类模式**:装饰器适合一次性逻辑(写一次、一个场景用);类模式适合需要参数化、跨 Agent 复用的治理规则。 > ### 4.3 动态提示词:让同一个 Agent 适配不同场景 有些 Agent 需要同时处理多种风格的任务——比如同一个报表助手,既要写周报(专业干练),又要写事故复盘(严谨客观),还要回答日常使用问题(简洁友好)。 硬写一个又大又全的 System Prompt 效果很差。`@dynamic_prompt` 根据每次请求的**实际内容动态生成 System Prompt**: ```python from langchain.agents.middleware import ModelRequest, dynamic_prompt def latest_user_text(request: ModelRequest) -> str: """从消息列表中提取最近一条用户消息的文本。""" # 倒序遍历,找到最新的 human 类型消息 for message in reversed(request.messages): if message.type == "human": return str(message.content) return "" @dynamic_prompt def adaptive_prompt(request: ModelRequest) -> str: """根据用户请求的内容,动态选择 System Prompt 风格。""" user_text = latest_user_text(request) # 从 runtime.context 读取受众信息(由外部系统传入) ctx = request.runtime.context or {} audience = ctx.get("audience", "内部团队") # 分支 1:事故复盘 → 严谨结构化 if "事故" in user_text or "复盘" in user_text: return ( "你是一名资深 SRE 工程师,正在撰写线上事故复盘报告。\n" f"目标受众:{audience}。\n" "请按以下结构组织内容:\n" "1. 事故概述(时间、影响范围、持续时间)\n" "2. 根因分析(直接原因、间接原因)\n" "3. 处理过程(时间线)\n" "4. 改进措施(短期 + 长期)\n" "语气:客观、严谨;不甩锅,不带情绪。" ) # 分支 2:周报 → 简洁干练 if "周报" in user_text: return ( "你是一名技术团队的 PM,正在编写项目周报。\n" f"目标受众:{audience}。\n" "请按以下结构组织内容:\n" "1. 本周完成事项(按优先级排列)\n" "2. 关键进展(数据支撑)\n" "3. 遇到的问题(附解决思路)\n" "4. 下周计划\n" "要求:简洁、有数据、不写流水账。" ) # 分支 3:默认日常回复 return ( "你是报表导出平台的智能助手。\n" f"目标受众:{audience}。\n" "请用简洁清晰的语言回答用户问题,必要时给出操作指引。" ) agent = create_agent( model=model, middleware=[adaptive_prompt], context_schema=RunContext, ) # 两个不同场景的请求,会得到完全不同的 System Prompt questions = [ "帮我写一份本周的报表系统开发周报,涵盖导出接口优化和权限模块改造", "帮我对上周六报表导出服务宕机的事故做一份复盘报告", ] for question in questions: result = agent.invoke( {"messages": [{"role": "user", "content": question}]}, context={"audience": "研发管理层"}, ) print(result["messages"][-1].content) print("---") ``` `@dynamic_prompt`** 适用场景**: + 同一个 Agent 处理多种不同类型的任务 + 不同用户角色需要不同的回复风格 + Prompt 需要根据上下文实时调整 **不适用场景**: + 审计日志 → 用 `@before_model` + 模型重试 → 用 `@wrap_model_call` + 工具拦截 → 用 `@wrap_tool_call` --- ## 附录:Middleware 选型速查表 ### 按关注维度分类 | 维度 | Middleware | 一句话说明 | 关键参数 | | --- | --- | --- | --- | | 🔢 次数限制 | `ModelCallLimitMiddleware` | 限制模型调用次数 | `run_limit`, `thread_limit`, `exit_behavior` | | 🔢 次数限制 | `ToolCallLimitMiddleware` | 限制特定工具调用次数 | `tool_name`, `run_limit`, `exit_behavior` | | 🛡️ 安全审批 | `HumanInTheLoopMiddleware` | 关键操作暂停等人类确认 | `interrupt_on`, `allowed_decisions` | | 🔒 隐私脱敏 | `PIIMiddleware` | 敏感信息进入模型前脱敏 | `detector`, `strategy`, `apply_to_*` | | 🔄 故障重试 | `ToolRetryMiddleware` | 工具失败后自动重试 | `max_retries`, `retry_on`, `on_failure` | | 🔄 故障重试 | `ModelRetryMiddleware` | 模型调用失败后重试 | `max_retries`, `retry_on` | | ⬇️ 模型降级 | `ModelFallbackMiddleware` | 主模型不可用时切换备用 | `fallback_model` | | 📝 上下文压缩 | `SummarizationMiddleware` | 长对话自动摘要压缩 | `trigger`, `keep`, `model`, `summary_prompt` | | 🧹 上下文清理 | `ContextEditingMiddleware` | 清理旧工具结果 | `edits=[ClearToolUsesEdit(...)]` | | 🎯 动态提示 | `@dynamic_prompt` | 每次模型调用前动态生成 System Prompt | 函数返回字符串 | ### 按执行阶段分类 wrap_model_call --> LLM[大模型] --> after_model end subgraph 工具调用阶段 wrap_tool_call --> Tool[工具执行] end subgraph 结束阶段 after_agent end --> ![](https://cdn.nlark.com/yuque/__mermaid_v3/c39852109b8074334ebca56007a74298.svg) | 阶段 | 可用的 Middleware 类型 | | --- | --- | | Agent 启动前 | 自定义 `before_agent` | | 模型调用前 | `@before_model`、`PIIMiddleware`(`apply_to_input`) | | 模型调用中 | `@wrap_model_call`、`ModelRetryMiddleware`、`ModelFallbackMiddleware` | | 模型调用后 | `@after_model`、`ModelCallLimitMiddleware`(计数判定) | | 工具调用中 | `@wrap_tool_call`、`ToolRetryMiddleware`、`ToolCallLimitMiddleware`、`HumanInTheLoopMiddleware` | | Agent 结束后 | 自定义 `after_agent` | | 上下文治理 | `SummarizationMiddleware`、`ContextEditingMiddleware` | | 提示词 | `@dynamic_prompt` | ### 组合建议:一份"生产级" Middleware 栈 把前面各章的内容组合起来,一份典型的**生产环境 Agent 配置**大概长这样: ```python agent = create_agent( model=local_model, tools=tools, middleware=[ # —— 安全层(最先执行) —— PIIMiddleware("email", strategy="redact", apply_to_input=True), HumanInTheLoopMiddleware( interrupt_on={"create_export_task": {"allowed_decisions": ["approve", "reject"]}} ), # —— 可靠性层 —— ModelRetryMiddleware(max_retries=3, retry_on=(Exception,)), ModelFallbackMiddleware(cloud_model), ToolRetryMiddleware( tools=["fetch_export_history"], max_retries=2, retry_on=(TimeoutError,) ), # —— 限流层 —— ModelCallLimitMiddleware(run_limit=20, exit_behavior="end"), ToolCallLimitMiddleware(tool_name="check_export_progress", run_limit=10), # —— 上下文治理层 —— SummarizationMiddleware(model=summary_model, trigger=("messages", 12), keep=("messages", 4)), ContextEditingMiddleware(edits=[ClearToolUsesEdit(trigger=3000, keep=2)]), # —— 自定义业务规则 —— ExportGovernanceMiddleware(role_limits={...}), # —— 动态提示词(最后,覆盖 System Prompt) —— adaptive_prompt, ], checkpointer=SqliteSaver.from_conn_string("checkpoints.db"), state_schema=ExportAgentState, context_schema=RunContext, system_prompt="你是报表导出平台的智能助手。", # 兜底 prompt,会被 @dynamic_prompt 覆盖 ) ``` > **Middleware 顺序很重要**。上面的栈里,脱敏在最前面(确保脏数据不进入后续环节),动态提示词在最后面(覆盖最终的 System Prompt)。调换顺序可能导致行为不一致,上线前务必在测试环境验证整条 Middleware 链的执行顺序。 > ---