LangChain 中间件实战指南:从 Python 异步到 MCP 集成,再到智能体治理

摘要:LangChain 的中间件体系是构建生产级智能体不可或缺的治理能力。本文系统梳理了 LangChain 第10章完整内容,从 Python 同步与异步的核心概念(async/await + asyncio)、MCP(Model Context Protocol)协议的实战集成(数学工具 + 天气 API),到 LangChain 1.0 预置中间件(Summarization、HITL)以及自定义中间件(基于装饰器 vs 基于类)的完整实现。每个知识点均配有可直接运行的代码示例,帮你快速掌握从基础到进阶的智能体治理技术。


目录

  • 一、Python 同步与异步:智能体提速的核心基础
    • 1.1 同步 vs 异步的本质区别
    • 1.2 Python 异步三要素:async / await / asyncio
    • 1.3 同步与异步的适用场景
    • 1.4 常见误区澄清
  • 二、MCP 的使用:让智能体调用"外部工具集"
    • 2.1 什么是 MCP
    • 2.2 安装依赖
    • 2.3 创建数学工具 MCP 服务端(stdio 模式)
    • 2.4 创建天气 MCP 服务端(HTTP 模式)
    • 2.5 创建 MCP 客户端并集成智能体
    • 2.6 整体架构梳理
  • 三、智能体中的预置中间件:开箱即用的治理能力
    • 3.1 LangChain 预置中间件全景
    • 3.2 Summarization(对话摘要)
    • 3.3 Human-in-the-loop(HITL,人工协作)
  • 四、自定义中间件:基于装饰器
    • 4.1 装饰器钩子类型
    • 4.2 完整实战:日志 + 重试 + 安全拦截 + 动态 Prompt
  • 五、自定义中间件:基于类
    • 5.1 类中间件的执行顺序规则
    • 5.2 完整实战:重试计时 + 安全拦截
  • 六、踩坑清单
  • 七、总结
  • 封面图提示词

一、Python 同步与异步:智能体提速的核心基础

1.1 同步 vs 异步的本质区别

你是否遇到过这样的场景:用 Python 写了一个下载图片的脚本,它却要一张下载完才能开始下一张,效率低得让人着急;或者写了个简单的 Web 服务,一个慢请求就能卡住后面所有用户?

其实,这背后藏着 Python 中同步与异步的核心逻辑——它们决定了程序如何"安排时间",直接影响着代码的效率。

同步:“排队办事,等完一个再一个”

想象你去银行办理业务,只有一个窗口(对应程序的"单线程"):

  • 你先取号排队,前面的人不办完,你只能等着
  • 前面的人可能要填单子、核对信息(对应程序中的"IO 操作"),哪怕他在填单子时窗口空闲,你也不能上前
  • 只有前一个人完全办完离开,下一个人才能开始

这就是同步编程的逻辑:任务按顺序执行,一个任务没完成(尤其是等待 IO 时),后面的任务必须"阻塞"(等待),直到当前任务结束。

同步代码示例

import time

# 模拟下载图片(IO操作,用time.sleep模拟等待时间)
def download_img(img_name):
    print(f"开始下载:{img_name}")
    time.sleep(2)  # 模拟下载等待,此时程序完全阻塞
    print(f"下载完成:{img_name}")

# 同步执行
start_time = time.time()
download_img("风景.jpg")
download_img("人物.jpg")
download_img("动物.jpg")
end_time = time.time()

print(f"总耗时:{end_time - start_time:.2f}秒")

运行结果:3 个任务依次执行,总耗时约 6 秒——这就是同步的"低效"之处:等待 IO 时,程序完全闲着

异步:“多件事穿插做,不等完也能切换”

还是刚才的银行场景,但这次你多了个"助手"(对应 Python 的"事件循环"):

  • 你去窗口提交申请后,工作人员说"填完单子再来"(对应程序发起 IO 请求)
  • 你不用在窗口等,而是去旁边填单子(程序释放线程,去处理其他任务
  • 同时,助手帮你盯着窗口——等你填完单子(IO 完成),助手会提醒你"该你了"(程序切换回原任务继续执行)
  • 这样一来,窗口和你都没闲着,效率大大提升

这就是异步编程的逻辑:任务发起后,如果需要等待 IO(比如下载、数据库查询),程序不会阻塞,而是去执行其他任务;等 IO 完成后,再回来继续处理之前的任务。

注意:异步不是"同时执行多个任务"(那是多线程/多进程),而是"在等待的间隙,穿插执行其他任务"——本质是"单线程下的任务调度优化"。

1.2 Python 异步三要素:async / await / asyncio

要素一:async def 定义异步函数

首先,异步函数必须用 async def 定义,而不是普通的 def——这告诉 Python:“这个函数是异步的,里面可能有等待操作”。

# 异步函数(用 async def 定义)
async def async_download_img(img_name):
    print(f"开始下载:{img_name}")
    # 这里不能用 time.sleep(同步等待),必须用异步等待
    await asyncio.sleep(2)  # 异步等待:释放线程去处理其他任务
    print(f"下载完成:{img_name}")

注意:普通函数里不能用 await,只有 async def 定义的异步函数里才能用。

要素二:await 等待"可等待对象"

await 的作用是"暂停当前异步函数,去执行其他任务,直到等待的对象完成"。它后面必须跟"可等待对象"(比如异步函数、asyncio.Task 等)。

为什么不能用 time.sleep(2)?

因为 time.sleep 是同步等待,会阻塞整个线程——而异步需要的是"异步等待",所以要用 asyncio.sleep(2)(它是异步的,会释放线程)。

要素三:asyncio 管理异步任务

asyncio 是 Python 标准库中专门用于异步编程的库,它提供了"事件循环"(相当于之前说的"助手")、任务调度、IO 管理等核心功能。

最常用的方法是 asyncio.run()——它会创建一个事件循环,运行异步函数,然后关闭循环

import asyncio
import time


async def download_img(image_name):
    print(f"开始下载图片:{image_name}")
    await asyncio.sleep(2)
    print("图片下载完成")


async def main():
    start_time = time.time()
    await asyncio.gather(
        download_img("风景.jpg"),
        download_img("人物.jpg"),
        download_img("宠物.jpg")
    )
    end_time = time.time()
    print(f"总耗时:{end_time - start_time}")


if __name__ == '__main__':
    asyncio.run(main())

运行结果:三个任务并发执行,总耗时约 2 秒——异步的魔力:三个任务的"等待时间"重叠了。

1.3 同步与异步的适用场景

维度 同步(Synchronous) 异步(Asynchronous)
执行方式 串行,逐个执行 穿插执行,等待时切到其他任务
适用场景 CPU 密集型、计算密集型 IO 密集型、网络请求、文件读写
典型例子 计算质数、矩阵运算 爬取网页、下载文件、数据库查询
优点 简单直观,易于调试 高并发,资源利用率高
缺点 IO 等待时资源浪费 需要学习成本,调试复杂

举个实际的例子

  • 如果你写一个"计算 1 到 100000 的质数"的程序(CPU 密集),用异步没用——因为 CPU 一直在工作,没有等待间隙,异步无法切换任务
  • 如果你写一个"爬取 100 个网页"的程序(IO 密集),用异步就很合适——因为爬网页时大部分时间在等服务器响应,异步可以在等待时爬其他网页

1.4 常见误区澄清

误区一:异步等于多线程/多进程

异步是"单线程下的任务调度",而多线程/多进程是"真正的并行执行"。异步不能解决 CPU 密集型任务的效率问题,因为它本质还是单线程。

误区二:不是所有库都支持异步

很多 Python 库是同步的(比如 requestssqlite3),在异步函数里用这些库会阻塞线程——必须用对应的异步库,比如 aiohttp(对应 requests)、asyncpg(对应 psycopg2)。

误区三:异步代码不一定比同步快

如果任务量很小,或者 IO 等待时间极短(比如本地文件读写),异步的"调度开销"可能比节省的等待时间还多,反而更慢。

总结:异步就是让 CPU 不要闲下来,在等待 IO 的间隙去干别的事。


二、MCP 的使用:让智能体调用"外部工具集"

2.1 什么是 MCP

MCP(Model Context Protocol,模型上下文协议)由 Anthropic 于 2024 年 11 月 25 日正式发布,用于标准化模型与外部工具、数据源和上下文环境之间的交互方式。

它的目标是让不同系统之间共享统一的模型上下文接口,使智能体能在不同的运行环境中调用工具时,无须关心底层传输细节。MCP 就像智能体世界的 “USB 协议”——无论工具在本地还是异地,使用 Python、NodeJS 或 HTTP,都可以通过统一接口被智能体调用。

MCP 本质上是一个客户端-服务端协议(Client-Server 协议):

  • MCP 服务端:负责定义并向网络暴露可用的工具
  • MCP 客户端:运行在智能体端,负责发现、加载并调用这些远程工具

LangChain 官方提供了 langchain-mcp-adapters 适配库来支持 MCP。

2.2 安装依赖

官方推荐使用 mcp 库来快速构建 MCP 服务端:

pip install mcp
pip install fastmcp

2.3 创建数学工具 MCP 服务端(stdio 模式)

创建一个数学工具 MCP 服务端 math_mcp_server.py

from mcp.server.fastmcp import FastMCP


mcp = FastMCP("Math")

# 定义第一个工具:加法
# 使用 @mcp.tool() 装饰器即可将函数注册为 MCP 工具。
@mcp.tool()
def add(a: int, b: int) -> int:
    """Add two numbers"""
    return a + b

# 定义第二个工具:乘法
# 使用 @mcp.tool() 装饰器即可将函数注册为 MCP 工具。
@mcp.tool()
def multiply(a: int, b: int) -> int:
    """Multiply two numbers"""
    return a * b

# 运行 MCP 服务,使用 stdio 传输协议
# transport="stdio" 表示通过标准输入输出进行通信。
if __name__ == "__main__":
    mcp.run(transport="stdio")

关键点

  • FastMCP("Math") 创建 MCP 服务端实例
  • @mcp.tool() 装饰器将函数注册为 MCP 工具
  • transport="stdio" 表示通过标准输入输出进行通信(本地子进程模式)

2.4 创建天气 MCP 服务端(HTTP 模式)

基于同样的方式,创建一个天气 MCP 服务端,调用第三方免费的 Open-Meteo API,并采用 streamable-http 模式启动:

from mcp.server.fastmcp import FastMCP
import httpx
from typing import Dict, Any, Optional

mcp = FastMCP("Weather")

# 定义 Open-Meteo 公共 API 地址(无需 API Key)
OPEN_METEO_WEATHER_URL = "https://api.open-meteo.com/v1/forecast"
OPEN_METEO_GEOCODE_URL = "https://geocoding-api.open-meteo.com/v1/search"

# 工具一:geocode_city()
# 将城市名解析为经纬度信息
@mcp.tool()
def geocode_city(name: str, country: Optional[str] = None, language: str = "zh") -> Dict[str, Any]:
    """将城市名解析为经纬度"""
    params = {"name": name, "count": 1, "language": language, "format": "json"}
    if country:
        params["country"] = country
    with httpx.Client(timeout=10) as client:
        r = client.get(OPEN_METEO_GEOCODE_URL, params=params)
        r.raise_for_status()
        data = r.json()
    results = data.get("results") or []
    if not results:
        return {"error": f"未找到城市:{name}"}
    top = results[0]
    return {
        "name": top.get("name"),
        "lat": top.get("latitude"),
        "lon": top.get("longitude"),
        "country": top.get("country"),
    }

# 工具二:get_current_weather()
# 根据经纬度查询当前天气
@mcp.tool()
def get_current_weather(lat: float, lon: float) -> Dict[str, Any]:
    """根据经纬度查询当前天气"""
    params = {"latitude": lat, "longitude": lon, "current_weather": True}
    with httpx.Client(timeout=10) as client:
        r = client.get(OPEN_METEO_WEATHER_URL, params=params)
        r.raise_for_status()
        payload = r.json()
    cw = payload.get("current_weather") or {}
    return {
        "latitude": lat,
        "longitude": lon,
        "temperature": cw.get("temperature"),
        "windspeed": cw.get("windspeed"),
        "weathercode": cw.get("weathercode"),
        "time": cw.get("time"),
    }

# 工具三:get_current_weather_by_city()
# 组合前两个工具,实现「城市名 → 当前天气」的完整流程
@mcp.tool()
def get_current_weather_by_city(name: str, country: Optional[str] = None, language: str = "zh") -> Dict[str, Any]:
    """城市名 -> 当前天气(内部先地理编码再查询天气)"""
    g = geocode_city(name=name, country=country, language=language)
    if "error" in g:
        return g
    w = get_current_weather(lat=g["lat"], lon=g["lon"])
    return {**g, **w}

# transport="streamable-http" 表示使用 HTTP 协议暴露服务
# 端点默认路径是 /mcp,端口通常为 8000
if __name__ == "__main__":
    print("Starting Weather MCP Server (streamable-http) on http://localhost:8000/mcp ...")
    mcp.run(transport="streamable-http")

创建完成后,通过以下命令启动该服务端

python weather_mcp_server.py

知识点补充:**g 展开字典传参

g = {"name": "张三", "age": 20}

def func(name, age):
    print(name, age)

# 等价写法1:手动传关键字
func(name="张三", age=20)

# 等价写法2:**g 自动拆解字典
func(**g)

2.5 创建 MCP 客户端并集成智能体

在编写 MCP 客户端代码之前,需要安装官方的 MCP 适配器依赖:

pip install langchain-mcp-adapters

随后,创建 MCP 客户端代码,该客户端将同时连接上述两个 MCP 服务端,并将获取到的远程工具注册到智能体中:

import asyncio

from langchain.agents import create_agent
from langchain_core.messages import HumanMessage
from langchain_mcp_adapters.client import MultiServerMCPClient

from utils.model_factory import get_deepSeek_model

model = get_deepSeek_model()

# 2. 连接 MCP 服务端 (Math + Weather)
client = MultiServerMCPClient({
    "Math": {
        "transport": "stdio",
        "command": "python",
        "args": ["./math_mcp_server.py"],  # 确保路径正确
    },
    "Weather": {
        "transport": "streamable_http",
        "url": "http://localhost:8000/mcp",  # Weather MCP Server
    },
})

tools = asyncio.run(client.get_tools())

agent = create_agent(
    model=model,
    tools=tools,
    system_prompt="你是一个助理。涉及数学计算,使用 Math 工具(add / multiply);"
        "涉及天气,使用 Weather 工具(geocode_city / get_current_weather / get_current_weather_by_city)。"
)

result1 = asyncio.run(agent.ainvoke({
    "messages": [
        {
            "role": "user",
            "content": "请帮我计算 (3 + 5) × 12 的结果"
        }
    ]
}))
for message in result1['messages']:
    print(type(message).__name__)
    print(message.content)


result2 = asyncio.run(agent.ainvoke({
    "messages": [
        HumanMessage(content="今天郑州天气如何?")
    ]
}))
for message in result2['messages']:
    print(type(message).__name__)
    print(message.content)

注意agent.invoke 本身是一种同步调用的方法,会阻塞线程。在异步环境中应使用 ainvoke

知识点补充:同步封装异步协程

import asyncio

def arun(coro):
    """同步封装: 把异步协程在顶层跑完,主逻辑仍然是「同步写法」"""
    return asyncio.run(coro)

async def demo():
    await asyncio.sleep(0.5)
    print("异步执行完成")
    return 666

# 主线程同步调用,不用 async/await
result = arun(demo())
print(result)  # 输出 666

运行输出

已加载的工具: ['add', 'multiply', 'geocode_city', 'get_current_weather', 'get_current_weather_by_city']

数学任务: 请帮我计算 (3 + 5) × 12 的结果

智能体输出: **(3 + 5) × 12 = 96**

计算过程:
1. 先算括号内:3 + 5 = 8
2. 再乘以 12:8 × 12 = 96

最终结果为 **96**。

天气任务: 请告诉我北京现在的天气情况

智能体输出: 以下是北京现在的天气情况:

- **城市**:北京(中国)
- **当前温度**:**27.6°C**
- **风速**:**11.0 km/h**
- **天气状况**:天气代码为 3(多云/阴天)
- **时间**:2026年6月22日 15:00

总体来说,北京现在气温约 **27.6°C**,体感比较温暖,风速适中,是多云或阴天的天气。

2.6 整体架构梳理

上述示例构建了一个完整的 MCP 集成环境,其架构与数据流如下:

组件 传输协议 暴露工具 说明
math_mcp_server.py stdio(标准输入输出) add、multiply 被客户端作为本地子进程拉起并通信
weather_mcp_server.py streamable-http geocode_city、get_current_weather、get_current_weather_by_city 基于 FastMCP 封装 Open-Meteo 公共 API,暴露标准 MCP HTTP 端点
client_mcp-demo.py 使用 MultiServerMCPClient 同时连接两个服务端,通过 client.get_tools 拉取工具定义,交给 create_agent 注册,智能体采用 ReAct 模式自动选择工具

三、智能体中的预置中间件:开箱即用的治理能力

为了提升智能体在生产环境中的项目化能力,LangChain 提供了一系列开箱即用的预置中间件(Prebuilt Middleware)。这些中间件从可靠性、可观测性、成本控制与安全治理等多个维度,为智能体系统提供了关键保障。

3.1 LangChain 预置中间件全景

中间件 作用
Summarization 在上下文长度接近上限时触发,压缩早期历史并保留最近消息,有效控制上下文长度与 API 成本
Human-in-the-loop(HITL) 对敏感或高风险的工具调用进行人工审批,支持批准、修改或拒绝执行,必要时中断流程等待人工决策
Anthropic prompt caching 对重复的 Prompt 段落进行缓存,避免重复计算,减少 Token 消耗
Model call limit 限制模型调用次数与频率,防止成本失控或陷入无限循环
Tool call limit 限制工具调用次数与频率,保护外部服务,保障系统稳定性与成本可控
Model fallback 在主模型调用失败或输出质量不达标时,自动切换至备用模型,提升系统整体鲁棒性
PII detection 在输入模型前或输出给用户前,自动检测并脱敏个人可识别信息,满足数据合规要求
To-do list 将复杂的多步任务拆解为结构化的待办事项清单,指导后续的工具调用与执行顺序
LLM tool selector 在调用主模型进行完整推理前,使用一个轻量级大模型预先智能筛选相关工具,提升决策效率
Tool retry 为工具调用提供可配置的重试与退避策略,增强对临时性失败的容错能力
LLM tool emulator 在缺乏真实工具环境或进行离线测试时,使用大模型模拟工具行为
Context editing 通过修剪、总结或清除历史工具调用记录等方式,主动管理对话上下文

注意:LangChain 预置的中间件仍在持续更新中,不同版本间可能在命名、参数或功能上有部分调整,请以官方文档为准。

3.2 Summarization(对话摘要)

用途:自动对对话历史进行压缩,避免超出模型的上下文窗口限制,同时尽力维持对话的连贯性。

关键机制:在模型调用前,检查上下文 Token 数量,若超过预设阈值则触发摘要进程,将早期消息压缩为简洁的摘要(以前聊天的大概核心内容),并可配置保留最近 N 条原始消息以维持细节

from langchain_openai import ChatOpenAI
from langchain.agents import create_agent
from langchain.agents.middleware import SummarizationMiddleware
from langchain.messages import HumanMessage

# 主对话模型
llm = ChatOpenAI(
    model="deepseek-chat",
    api_key="sk-f3c7557e9fc54541802fd638de03bbeb",
    base_url="https://api.deepseek.com",
    temperature=0.0,
    max_tokens=20,  # 限制每次生成不超过 20 tokens
)

# 摘要模型(这里也使用 deepseek)
summ_llm = ChatOpenAI(
    model="deepseek-chat",
    api_key="sk-f3c7557e9fc54541802fd638de03bbeb",
    base_url="https://api.deepseek.com",
    temperature=0.0,
    max_tokens=20,  # 限制每次生成不超过 20 tokens
    model_provider="openai"
)

# 会话摘要中间件:极低阈值,快速触发;摘要后仅保留 1 条原文
middleware = SummarizationMiddleware(
    model=summ_llm,
    max_tokens_before_summary=1000,  # 达到 1000 token 触发摘要
    messages_to_keep=2,  # 保留最近 2 条消息
    summary_prompt="用20个字以内概括要点。",  # 自定义"如何摘要"的提示词
)

agent = create_agent(
    model=llm,
    tools=[],
    system_prompt="只用一句极短中文回答,且不超过30个字。",  # 只用"极简短句"回答
    middleware=[middleware],
)


conversation = [
    "RAG是什么?",
    "有何用途?",
    "主要缺点?",
    "一句话总结"
]

state = {"messages": []}

for i, question in enumerate(conversation, 1):
    # 用户问题
    user_msg = HumanMessage(content=question)
    # 智能体响应
    state = agent.invoke({"messages": state["messages"] + [user_msg]})
    answer = state["messages"][-1].content

    print(f"\n第 {i} 轮")
    print(f"Q: {question}")
    print(f"A: {answer}")

print("\n对话结束。")

SummarizationMiddleware 关键参数说明

参数 说明
max_tokens_before_summary 触发自动摘要的 Token 阈值。对话累计 Token 临近/超过该值,中间件把早期消息压缩为摘要,控制上下文长度与调用成本
messages_to_keep 摘要完成后,保留最近 N 条原始消息不压缩,防止全量摘要丢失细节
summary_prompt(可选) 自定义摘要提示词,控制摘要风格、粒度(精简 / 保留关键数据);不配置则使用框架默认摘要 Prompt

运行结果

第 1 轮
Q: RAG是什么?
A: RAG是检索增强生成技术。

第 2 轮
Q: 有何用途?
A: 提升回答准确性和知识广度。

第 3 轮
Q: 主要缺点?
A: 依赖检索质量,可能增加延迟。

第 4 轮
Q: 一句话总结
A: RAG通过检索增强生成,提升准确但依赖检索质量。

对话结束。

结果原理说明

第 4 轮回答看似重复第 3 轮,非模型错误。示例阈值 30 Token 配置偏低,第 3 轮对话结束后上下文总量触达阈值,第 4 轮模型调用前中间件执行摘要:早期历史被压缩、丢失 RAG 定义/用途/缺点细节,上下文只剩摘要 + 就近留存消息;模型仅依托精简后的上下文推理,因此沿用最近内容作答。

3.3 Human-in-the-loop(HITL,人工协作)

用途:对标记为高风险的工具调用实施人工审批流程。支持批准、编辑参数后执行或拒绝操作,并在需要时中断智能体执行流程,等待人工输入。

关键机制:在工具调用前的生命周期钩子中进行策略检查,若命中规则则暂停执行并抛出中断信号,等待外部系统或用户做出决策。

from langchain.chat_models import init_chat_model
from langchain_openai import ChatOpenAI
from langchain.agents import create_agent
from langchain.agents.middleware import HumanInTheLoopMiddleware
from langchain.tools import tool
from langchain.messages import HumanMessage
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import Command

# 1) 主模型
llm = init_chat_model(
    model="deepseek-chat",
    api_key="sk-f3c7557e9fc54541802fd638de03bbeb",
    base_url="https://api.deepseek.com",
    temperature=0.0,
    max_tokens=20,
    model_provider="openai"
)

# 2) 高风险工具(示例:数据库写操作)
@tool
def dangerous_write(sql: str) -> str:
    """对数据库执行写操作(插入/更新/删除)。当请求涉及数据库写操作或出现以'SQL:'开头的指令时,必须调用本工具。"""
    return f"[模拟执行] {sql}"

# 3) HITL 中间件:拦截指定工具
hitl = HumanInTheLoopMiddleware(
    interrupt_on={"dangerous_write": {"allowed_decisions": ["approve", "edit", "reject"]}}
)

# 4) 创建 Agent(强提示:遇到 SQL 必须用 dangerous_write)
agent = create_agent(
    model=llm,
    tools=[dangerous_write],
    middleware=[hitl],
    system_prompt=(
        "只用一句极短中文回答(≤20字)。"
        "凡是涉及数据库写操作,或消息以SQL:开头时,必须调用工具 dangerous_write,不得直接回答。"
    ),
    checkpointer=InMemorySaver(),  # HITL 必须
)

CFG = {"configurable": {"thread_id": "hitl-demo-interactive"}}

# 5) 四轮对话:第2轮与第4轮都触发 HITL;让你交互输入决定
conversation = [
    "你是谁?",  # 安全问答
    "SQL: INSERT INTO logs(content) VALUES ('hello');",  # 触发 HITL(预计输入 approve)
    "继续。",  # 安全问答
    "SQL: DELETE FROM orders WHERE created_at >= date('now','-7 days');",  # 触发 HITL(预计输入 reject)
]

state = {"messages": []}

def handle_interrupt(result) -> dict:
    """处理 HITL 中断:从命令行读取决策,并用 Command(resume=...) 恢复。"""
    interrupt = result.get("__interrupt__")
    if not interrupt:
        return result

    print("\n检测到人工在环中断:")
    print(interrupt)  # 打印被拦截的工具与参数

    decision = input("请输入决策 (approve/edit/reject):").strip().lower()

    if decision == "approve":
        # 直接放行
        return agent.invoke(Command(resume={"decisions": [{"type": "approve"}]}), config=CFG)

    elif decision == "edit":
        # 允许修改参数(例如 SQL),再继续
        new_sql = input("请输入修改后的 SQL:").strip()
        return agent.invoke(
            Command(
                resume={
                    "decisions": [{
                        "type": "edit",
                        "edited_action": {
                            "name": "dangerous_write",
                            "args": {"sql": new_sql}
                        }
                    }]
                }
            ),
            config=CFG
        )

    elif decision == "reject":
        # 拒绝执行,并给出替代返回(不执行工具)
        return agent.invoke(
            Command(
                resume={
                    "decisions": [{
                        "type": "reject",
                        "override": {"content": "[操作已被人工拒绝]"}
                    }]
                }
            ),
            config=CFG
        )
    else:
        print("输入无效,按拒绝处理。")
        return agent.invoke(
            Command(
                resume={
                    "decisions": [{
                        "type": "reject",
                        "override": {"content": "[操作已被人工拒绝]"}
                    }]
                }
            ),
            config=CFG
        )

for i, q in enumerate(conversation, 1):
    user_msg = HumanMessage(content=q)
    result = agent.invoke({"messages": state["messages"] + [user_msg]}, config=CFG)

    # 命中中断:让你做决策 → 再 resume
    if result.get("__interrupt__"):
        result = handle_interrupt(result)

    state = result
    ans = state["messages"][-1].content

    print(f"\n第 {i} 轮")
    print(f"Q: {q}")
    print(f"A: {ans}")

print("\n对话结束。")

核心机制解析

  1. 当智能体尝试调用需审批的工具(如 dangerous_write)时,该中间件会自动中断执行流程,等待人工确认
  2. 通过 InMemorySaver 保存执行状态,确保后续可以通过相同的 thread_id 从断点恢复执行
  3. 人工决策支持三种模式:
    • approve:批准并继续执行工具调用
    • edit:修改工具调用的输入参数后继续执行
    • reject:拒绝执行该工具,并向智能体返回一个替代性的提示信息

四、自定义中间件:基于装饰器

在熟悉了 LangChain 预置中间件后,我们开始探索如何创建自定义中间件。自定义中间件通过在智能体执行流程的特定钩子(hook)上注册逻辑来实现。

4.1 装饰器钩子类型

LangChain 官方将可用的钩子分为三类:

类型 钩子 触发时机 典型用途
节点式(Node-style) @before_agent 智能体执行前 初始化、上下文检查、权限验证
节点式(Node-style) @before_model 每次模型调用前 修改输入、记录日志、执行过滤
节点式(Node-style) @after_model 每次模型响应后 分析输出、检测风险词、修正结果
节点式(Node-style) @after_agent 智能体执行完成后 日志上报、状态持久化
包裹式(Wrap-style) @wrap_model_call 每次模型调用时 计时、限流、重试、异常兜底
包裹式(Wrap-style) @wrap_tool_call 每次工具调用时 审计、权限校验、记录调用轨迹
便捷装饰器 @dynamic_prompt 模型调用前 动态修改系统 Prompt

4.2 完整实战:日志 + 重试 + 安全拦截 + 动态 Prompt

from typing import Any, Callable
import time
import math

from langchain.agents import create_agent
from langchain.agents.middleware import (
    before_agent,
    after_agent,
    before_model,
    after_model,
    wrap_model_call,
    dynamic_prompt,
    AgentState,
    ModelRequest,
    ModelResponse,
)
from langchain.messages import AIMessage, HumanMessage
from langgraph.runtime import Runtime
from langchain_openai import ChatOpenAI


# ========= 便捷装饰器:动态系统提示 =========
@dynamic_prompt
def personalized_prompt(req: ModelRequest) -> str:
    runtime = getattr(req, "runtime", None)
    user_id = "访客"
    if runtime and getattr(runtime, "context", None):
        user_id = runtime.context.get("user_id", "访客")
    return f"你是一名贴心的中文助手,正在为用户「{user_id}」提供帮助。回答时要简洁、自然。"


# ========= 节点式:Agent 级别前置/后置 =========
@before_agent
def log_before_agent(state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
    print(f"[before_agent] 本次会话开始,已有消息数:{len(state.get('messages', []))}")
    return None


@after_agent
def log_after_agent(state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
    print(f"[after_agent] 会话结束,最终消息数:{len(state.get('messages', []))}")
    return None


# ========= 节点式:模型调用前/后 =========
@before_model
def log_before_model(state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
    print(f"[before_model] 准备进行模型调用,当前消息数:{len(state['messages'])}")
    return None


@after_model(can_jump_to=["end"])
def validate_output(state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
    """简单的输出校验:若模型输出包含 BLOCKED_CN,则改写消息并跳转到 end。"""
    last = state["messages"][-1]
    if isinstance(last, AIMessage) and "BLOCKED_CN" in (last.content or ""):
        print("[after_model] 触发安全规则:检测到 BLOCKED_CN,跳转到 end")
        return {
            "messages": [AIMessage("该请求触发了安全校验,无法继续。")],
            "jump_to": "end",
        }
    return None


# ========= 包裹式:为模型调用加重试与耗时统计 =========
@wrap_model_call
def retry_and_timing(
    request: ModelRequest,
    handler: Callable[[ModelRequest], ModelResponse],
) -> ModelResponse:
    max_retries = 2
    start = time.time()
    try:
        for i in range(max_retries + 1):
            try:
                return handler(request)
            except Exception as e:
                if i == max_retries:
                    raise
                # 退避等待(100ms, 200ms)
                backoff = 0.1 * math.pow(2, i)
                print(f"[wrap_model_call] 调用失败,将在 {backoff:.2f}s 后重试:{e}")
                time.sleep(backoff)
    finally:
        cost = (time.time() - start) * 1000
        print(f"[wrap_model_call] 本次模型调用耗时:{cost:.0f} ms")


# ========= 组装 Agent(DeepSeek 模型) =========
llm = ChatOpenAI(
    model="deepseek-chat",
    api_key="sk-f3c7557e9fc54541802fd638de03bbeb",
    base_url="https://api.deepseek.com",
    temperature=0.3,
)

agent = create_agent(
    model=llm,
    tools=[],  # 如需可加工具,此处保持最小化
    middleware=[
        personalized_prompt,  # 便捷:动态系统提示
        log_before_agent,     # 节点:Agent 级前置
        log_before_model,     # 节点:模型前
        retry_and_timing,     # 包裹:重试与计时
        validate_output,      # 节点:模型后(含 jump_to)
        log_after_agent,      # 节点:Agent 级后置
    ],
)

# ========= 最小演示 =========
if __name__ == "__main__":
    # 1) 正常问答(不会触发安全跳转)
    res1 = agent.invoke(
        {"messages": [HumanMessage("用一句话解释 LangGraph 是什么。")]},
        config={"context": {"user_id": "alice"}},
    )
    print("\n[Result-1]", res1["messages"][-1].content)

    print("--------------------------------")

    # 2) 触发 after_model 的阻断(让模型输出包含关键字)
    res2 = agent.invoke(
        {"messages": [HumanMessage("请只回复:BLOCKED_CN")]},
        config={"context": {"user_id": "bob"}},
    )
    print("\n[Result-2]", res2["messages"][-1].content)

运行结果

[before_agent] 本次会话开始,已有消息数:1
[before_model] 准备进行模型调用,当前消息数:1
[wrap_model_call] 本次模型调用耗时:999 ms
[after_agent] 会话结束,最终消息数:2

[Result-1] LangGraph 是一个用于构建有状态、多步骤语言模型工作流的框架,支持图结构编排和动态控制。
--------------------------------
[before_agent] 本次会话开始,已有消息数:1
[before_model] 准备进行模型调用,当前消息数:1
[wrap_model_call] 本次模型调用耗时:880 ms
[after_model] 触发安全规则:检测到 BLOCKED_CN,跳转到 end
[after_agent] 会话结束,最终消息数:3

[Result-2] 该请求触发了安全校验,无法继续。

要点解析

  1. 第一次调用,模型正常返回,触发了完整的日志与耗时打印
  2. 第二次调用,模型输出包含关键字 BLOCKED_CN@after_model 成功识别并跳转到 end,返回了安全提示
  3. 注意第二次有 3 条消息:问题 HumanMessage → 模型原始回答 AIMessage → 被拦截后替换的新 AIMessage,消息是追加而非覆盖

五、自定义中间件:基于类

当中间件逻辑变得复杂,需要组合使用多个钩子,或需要维护内部状态以支持可配置性与复用时,基于类(Class-based)的中间件实现方式便成为更合适的选择。

5.1 类中间件的执行顺序规则

钩子类型 执行顺序 说明
before_* 钩子 按中间件在列表中的注册顺序执行(从左到右) 先注册的先执行
包裹式钩子 外层的中间件先于内层的中间件执行 形成"洋葱模型"
after_* 钩子 执行顺序与 before_* 相反(从右到左) 后注册的先执行

5.2 完整实战:重试计时 + 安全拦截

from typing import Any, Callable
import time, math

from langchain.agents import create_agent
from langchain.agents.middleware import (
    AgentMiddleware, AgentState, ModelRequest, ModelResponse,
)
from langchain.messages import AIMessage, HumanMessage
from langgraph.runtime import Runtime
from langchain_openai import ChatOpenAI

# ========== 中间件 A:日志 + 安全拦截(Node-style) ==========
class PolicyGuardMiddleware(AgentMiddleware):
    """在关键节点打印日志,并在 after_model 命中关键词时跳转 end"""

    def before_agent(self, state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
        print(f"[before_agent] 会话开始,消息数:{len(state.get('messages', []))}")
        return None

    def before_model(self, state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
        print(f"[before_model] 准备进行模型调用,当前消息数:{len(state['messages'])}")
        return None

    def after_model(self, state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
        last = state["messages"][-1]
        if isinstance(last, AIMessage) and "BLOCKED_CN" in (last.content or ""):
            print("[after_model] 触发安全规则:检测到 BLOCKED_CN,跳转 end")
            return {
                "messages": [AIMessage("该请求触发了安全校验,无法继续。")],
                "jump_to": "end",
            }
        return None

    def after_agent(self, state: AgentState, runtime: Runtime) -> dict[str, Any] | None:
        print(f"[after_agent] 会话结束,最终消息数:{len(state.get('messages', []))}")
        return None


# ========== 中间件 B:重试 + 耗时统计(Wrap-style) ==========
class RetryAndMetricsMiddleware(AgentMiddleware):
    """包裹每次模型调用,做退避重试与耗时统计"""

    def __init__(self, max_retries: int = 2, base_delay: float = 0.1):
        self.max_retries = max_retries
        self.base_delay = base_delay  # 秒

    def wrap_model_call(
        self,
        request: ModelRequest,
        handler: Callable[[ModelRequest], ModelResponse],
    ) -> ModelResponse:
        start = time.time()
        try:
            for i in range(self.max_retries + 1):
                try:
                    return handler(request)
                except Exception as e:
                    if i == self.max_retries:
                        raise
                    backoff = self.base_delay * math.pow(2, i)
                    print(f"[wrap_model_call] 失败,将在 {backoff:.2f}s 后重试:{e}")
                    time.sleep(backoff)
        finally:
            cost = (time.time() - start) * 1000
            print(f"[wrap_model_call] 本次模型调用耗时:{cost:.0f} ms")


# ========== 组装 Agent(DeepSeek 模型) ==========
llm = ChatOpenAI(
    model="deepseek-chat",
    api_key="sk-f3c7557e9fc54541802fd638de03bbeb",
    base_url="https://api.deepseek.com",
    temperature=0.3
)

# 注意执行顺序:
# - wrap 型按列表"从外到内"包裹;我们希望"重试/计时"最外层,所以把 RetryAndMetricsMiddleware 放前面
# - node 型按列表顺序执行 before_*,after_* 逆序回卷
agent = create_agent(
    model=llm,
    tools=[],
    middleware=[
        RetryAndMetricsMiddleware(max_retries=2, base_delay=0.1),  # 外层:重试+计时
        PolicyGuardMiddleware(),                                     # 内层:日志+安全拦截
    ],
)

if __name__ == "__main__":
    # 1) 正常问答
    res1 = agent.invoke({"messages": [HumanMessage("用一句话解释 LangGraph 是什么。")]})
    print("\n[Result-1]", res1["messages"][-1].content)

    print("\n" + "-" * 64 + "\n")

    # 2) 触发安全跳转(让模型回复包含关键字)
    res2 = agent.invoke({"messages": [HumanMessage("请只回复:BLOCKED_CN")]})
    print("\n[Result-2]", res2["messages"][-1].content)

运行结果

[before_agent] 会话开始,消息数:1
[before_model] 准备进行模型调用,当前消息数:1
[wrap_model_call] 本次模型调用耗时:1283 ms
[after_agent] 会话结束,最终消息数:2

[Result-1] LangGraph 是一个用于构建有状态、多步骤的 AI 代理工作流的框架,它通过有向图结构来编排大语言模型(LLM)的调用、工具使用和状态管理。

----------------------------------------------------------------

[before_agent] 会话开始,消息数:1
[before_model] 准备进行模型调用,当前消息数:1
[wrap_model_call] 本次模型调用耗时:680 ms
[after_model] 触发安全规则:检测到 BLOCKED_CN,跳转 end
[after_agent] 会话结束,最终消息数:3

[Result-2] 该请求触发了安全校验,无法继续。

代码结构与职责解析

  1. RetryAndMetricsMiddleware(Wrap-style):包裹每次模型调用,添加退避重试与耗时统计,像洋葱最外层,会在真实模型调用前后插入控制与度量逻辑
  2. PolicyGuardMiddleware(Node-style):在关键节点打印日志,在 after_model 检查敏感关键字 BLOCKED_CN,命中时改写输出并跳转到 end

为什么把"重试/计时"放在前面、把"日志/安全拦截"放在后面?

  • wrap 型按列表从左到右成为外→内的包裹层,RetryAndMetricsMiddleware 在最外层,先于一切模型调用执行
  • before_* 钩子按列表从左到右执行(先外后内)
  • after_* 钩子按列表从右到左执行(先内后外)
  • 因为 PolicyGuardMiddleware 在列表靠后,所以它的 after_model 会优先得到模型响应,从而更早地做安全拦截与 jump_to="end" 跳转

这就是外层先兜住稳定性(重试 & 计时),内层更早拿到真实响应做策略判断(拦截 / 跳转)的原因。


六、踩坑清单

描述 解决方案
同步库在异步函数中阻塞 time.sleeprequests 等同步库会阻塞事件循环 使用 asyncio.sleepaiohttp 等异步替代方案
代码块语言标识缺失 CSDN 编辑器要求代码块标注语言才能正确高亮 所有代码块首行必须标注语言(pythonplain 等)
HITL 缺少 checkpointer HumanInTheLoopMiddleware 必须配合 InMemorySaver() 使用 create_agent 中传入 checkpointer=InMemorySaver()
工具路径配置错误 MCP stdio 模式的 args 路径需要绝对或相对当前工作目录 确保 math_mcp_server.py 路径正确,或改为绝对路径
异步封装问题 agent.invoke 是同步方法,在异步环境调用会阻塞 使用 ainvokeasyncio.run() 正确封装
中间件执行顺序混乱 Wrap-style 和 Node-style 执行顺序规则不同 外层放稳定性中间件(重试、计时),内层放策略中间件(安全、拦截)
摘要阈值配置过低 max_tokens_before_summary 配置过低导致过早压缩丢失上下文 根据模型上下文窗口和对话长度合理配置阈值
旧版 API 不兼容 LangChain 1.0 中间件 API 与早期版本有 Breaking Change 使用 langchain.agents.middleware 新模块而非旧版 API
MCP 服务端端口冲突 streamable-http 默认使用 8000 端口 检查端口占用,或通过 mcp.run(port=8001) 指定其他端口
嵌套装饰器层数过多 多个装饰器叠加时,执行顺序可能不符合直觉 遵循"洋葱模型"理解包裹式钩子的内外层关系

七、总结

LangChain 的智能体与中间件共同构成了其架构中核心的能力体系。LangChain v1.0 对这一体系进行了彻底重构,确立了清晰的架构分层:智能体专注于决策执行,而中间件则承担过程治理

本文从三个层面帮你建立了完整的治理框架:

  1. Python 异步基础:掌握 async/awaitasyncio 是高效构建智能体的前提。异步的核心是让 CPU 在等待 IO 时不闲着,这是提升智能体吞吐量的关键。

  2. MCP 协议集成:通过 FastMCP + MultiServerMCPClient,你可以轻松将任何 Python 函数(数学运算、天气查询)封装为 MCP 工具,并通过 stdio 或 HTTP 传输协议暴露给智能体调用。MCP 就像智能体的"USB 协议",统一了工具调用的接口标准。

  3. 中间件治理能力

    • 预置中间件:Summarization 控制上下文长度和成本,HITL 保障高风险操作的人工审批
    • 自定义装饰器中间件:轻量、灵活,适合单一职责的场景(日志、计时、安全拦截)
    • 自定义类中间件:可配置、可复用、可维护内部状态,适合复杂的多钩子组合场景
中间件类型 适用场景 复杂度 推荐优先级
预置中间件 通用治理需求(摘要、人工审批) 首选
基于装饰器的中间件 单一职责的简单逻辑 次选
基于类的中间件 多钩子组合、状态维护、可配置复用 进阶

中间件跟智能体的业务没多少关系,它一般用于在智能体的某些生命周期中自动触发,可以帮我们进行权限校验、日志记录、消息压缩、人为介入等——有点像 Spring 中的 AOP。


如果觉得这篇文章对你有帮助,欢迎点赞收藏。你在使用 LangChain 中间件时遇到了什么问题?评论区交流。

Logo

欢迎加入 MCP 技术社区!与志同道合者携手前行,一同解锁 MCP 技术的无限可能!

更多推荐