# LangChain 中间件实战指南:从 Python 异步到 MCP 集成,再到智能体治理
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 库是同步的(比如 requests、sqlite3),在异步函数里用这些库会阻塞线程——必须用对应的异步库,比如 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对话结束。")
核心机制解析:
- 当智能体尝试调用需审批的工具(如
dangerous_write)时,该中间件会自动中断执行流程,等待人工确认 - 通过
InMemorySaver保存执行状态,确保后续可以通过相同的thread_id从断点恢复执行 - 人工决策支持三种模式:
- 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] 该请求触发了安全校验,无法继续。
要点解析:
- 第一次调用,模型正常返回,触发了完整的日志与耗时打印
- 第二次调用,模型输出包含关键字
BLOCKED_CN,@after_model成功识别并跳转到end,返回了安全提示 - 注意第二次有 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] 该请求触发了安全校验,无法继续。
代码结构与职责解析:
- RetryAndMetricsMiddleware(Wrap-style):包裹每次模型调用,添加退避重试与耗时统计,像洋葱最外层,会在真实模型调用前后插入控制与度量逻辑
- PolicyGuardMiddleware(Node-style):在关键节点打印日志,在
after_model检查敏感关键字BLOCKED_CN,命中时改写输出并跳转到end
为什么把"重试/计时"放在前面、把"日志/安全拦截"放在后面?
wrap型按列表从左到右成为外→内的包裹层,RetryAndMetricsMiddleware 在最外层,先于一切模型调用执行before_*钩子按列表从左到右执行(先外后内)after_*钩子按列表从右到左执行(先内后外)- 因为 PolicyGuardMiddleware 在列表靠后,所以它的
after_model会优先得到模型响应,从而更早地做安全拦截与jump_to="end"跳转
这就是外层先兜住稳定性(重试 & 计时),内层更早拿到真实响应做策略判断(拦截 / 跳转)的原因。
六、踩坑清单
| 坑 | 描述 | 解决方案 |
|---|---|---|
| 同步库在异步函数中阻塞 | time.sleep、requests 等同步库会阻塞事件循环 |
使用 asyncio.sleep、aiohttp 等异步替代方案 |
| 代码块语言标识缺失 | CSDN 编辑器要求代码块标注语言才能正确高亮 | 所有代码块首行必须标注语言(python、plain 等) |
| HITL 缺少 checkpointer | HumanInTheLoopMiddleware 必须配合 InMemorySaver() 使用 |
在 create_agent 中传入 checkpointer=InMemorySaver() |
| 工具路径配置错误 | MCP stdio 模式的 args 路径需要绝对或相对当前工作目录 |
确保 math_mcp_server.py 路径正确,或改为绝对路径 |
| 异步封装问题 | agent.invoke 是同步方法,在异步环境调用会阻塞 |
使用 ainvoke 或 asyncio.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 对这一体系进行了彻底重构,确立了清晰的架构分层:智能体专注于决策执行,而中间件则承担过程治理。
本文从三个层面帮你建立了完整的治理框架:
-
Python 异步基础:掌握
async/await和asyncio是高效构建智能体的前提。异步的核心是让 CPU 在等待 IO 时不闲着,这是提升智能体吞吐量的关键。 -
MCP 协议集成:通过 FastMCP + MultiServerMCPClient,你可以轻松将任何 Python 函数(数学运算、天气查询)封装为 MCP 工具,并通过 stdio 或 HTTP 传输协议暴露给智能体调用。MCP 就像智能体的"USB 协议",统一了工具调用的接口标准。
-
中间件治理能力:
- 预置中间件:Summarization 控制上下文长度和成本,HITL 保障高风险操作的人工审批
- 自定义装饰器中间件:轻量、灵活,适合单一职责的场景(日志、计时、安全拦截)
- 自定义类中间件:可配置、可复用、可维护内部状态,适合复杂的多钩子组合场景
| 中间件类型 | 适用场景 | 复杂度 | 推荐优先级 |
|---|---|---|---|
| 预置中间件 | 通用治理需求(摘要、人工审批) | 低 | 首选 |
| 基于装饰器的中间件 | 单一职责的简单逻辑 | 中 | 次选 |
| 基于类的中间件 | 多钩子组合、状态维护、可配置复用 | 高 | 进阶 |
中间件跟智能体的业务没多少关系,它一般用于在智能体的某些生命周期中自动触发,可以帮我们进行权限校验、日志记录、消息压缩、人为介入等——有点像 Spring 中的 AOP。
如果觉得这篇文章对你有帮助,欢迎点赞收藏。你在使用 LangChain 中间件时遇到了什么问题?评论区交流。
更多推荐



所有评论(0)