AI Agent还在裸奔?MCP+中间件+异步——给智能体穿上生产级铁甲的完整工程指南
免责声明一:本文基于 LangChain 1.0+ 版本编写。
create_agent是 LangChain 1.0 引入的 API,旧版本(0.3.x)请使用create_react_agent+AgentExecutor。不同版本间 API 可能有差异,请以官方文档为准。免责声明二:本文所有代码均为教学示例,可直接复制运行(需配置 API Key)。生产环境请根据实际需求调整参数、增加异常处理和日志记录。本文遵循 MIT 开源协议,可自由使用和传播。
密钥安全提醒:本文代码中所有 API Key 均通过
os.getenv()从环境变量读取,请勿将密钥硬编码在源码中或提交到 Git 仓库。建议使用.env文件管理敏感信息,并在.gitignore中排除。前置知识门槛:阅读本文前,建议先掌握以下内容——
- Python 基础语法(函数、类、装饰器)
- LangChain 基本概念(Model、Tool、Agent)
- 大模型 API 调用基础(如 DeepSeek/OpenAI API)
多模型兼容说明:本文代码以 DeepSeek(
deepseek-v4-flash)为例,但所有示例均兼容 OpenAI GPT-4o、阿里通义千问等主流大模型,只需更换model_factory.py中的模型实例化代码即可。
AI Agent 还在裸奔?MCP+中间件+异步——给智能体穿上生产级铁甲的完整工程指南
你遇到过这些问题吗?
Agent 写好了,Demo 跑得很溜,但一上生产环境就翻车——
- 高并发场景下 Agent 响应慢如蜗牛,一个请求卡住后面全排队?
- 想让 AI 调用外部工具(天气、数据库、第三方 API),但每接一个工具就写一堆胶水代码?
- Agent 调用大模型时没有日志、没有重试、没有限流,出了问题全靠猜?
- 敏感操作(删数据库、执行 SQL)没有人工审批,Agent 说干就干?
如果你中了任何一条,说明你的 Agent 还在裸奔——没有穿上生产级的铁甲。
今天这篇文章,我们就来把这副铁甲一片一片装上:
┌─────────────────────────────────────────────────────────────────┐
│ Agent 生产化三板斧 │
├──────────────┬──────────────────┬───────────────────────────────┤
│ 第一板斧 │ 第二板斧 │ 第三板斧 │
│ 异步编程 │ MCP 协议 │ 中间件 │
│ 解决并发效率 │ 解决工具标准化 │ 解决安全/可靠/可观测 │
├──────────────┼──────────────────┼───────────────────────────────┤
│ async/await │ FastMCP │ SummarizationMiddleware │
│ asyncio │ MultiServerMCPClient │ HumanInTheLoopMiddleware │
│ asyncio.gather│ stdio/http传输 │ @before_model / @after_model │
│ │ │ @wrap_model_call │
│ │ │ Class-based Middleware │
└──────────────┴──────────────────┴───────────────────────────────┘
Block 1:Python 异步编程——让程序「不闲着」
1.1 同步 vs 异步:到底在「同步」什么?
同步和异步的核心,是程序处理任务的顺序和等待逻辑。
同步——「排队办事,等完一个再一个」。
想象你去银行办理业务,只有一个窗口(对应程序的「单线程」):
- 你取号排队,前面的人不办完,你只能等着
- 前面的人可能要填单子、核对信息(对应程序中的 IO 操作,比如下载文件、访问数据库),哪怕他在填单子时窗口空闲,你也不能上前
- 只有前一个人完全办完离开,下一个人才能开始
用 Python 代码举一个同步的例子(模拟下载 3 张图片):
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}秒")
运行结果:
开始下载:风景.jpg
下载完成:风景.jpg
开始下载:人物.jpg
下载完成:人物.jpg
开始下载:动物.jpg
下载完成:动物.jpg
总耗时:6.00秒
3 个任务依次执行,总耗时是单个任务的 3 倍——这就是同步的「低效」之处:等待 IO 时,程序完全闲着。
异步——「多件事穿插做,不等完也能切换」。
还是刚才的银行场景,但这次你多了个「助手」(对应 Python 的「事件循环」):
- 你去窗口提交申请后,工作人员说「填完单子再来」(对应程序发起 IO 请求)
- 你不用在窗口等,而是去旁边填单子(程序释放线程,去处理其他任务)
- 同时,助手帮你盯着窗口——等你填完单子(IO 完成),助手会提醒你「该你了」(程序切换回原任务继续执行)
- 这样一来,窗口和你都没闲着,效率大大提升
这就是异步编程的逻辑:任务发起后,如果需要等待 IO(比如下载、数据库查询),程序不会阻塞,而是去执行其他任务;等 IO 完成后,再回来继续处理之前的任务。
注意:异步不是「同时执行多个任务」(那是多线程/多进程),而是「在等待的间隙,穿插执行其他任务」——本质是单线程下的任务调度优化。
1.2 异步三要素:async / await / asyncio
要在 Python 中写异步代码,有三个核心概念必须掌握。
要素一:async def 定义异步函数
异步函数必须用 async def 定义,而不是普通的 def——这告诉 Python:「这个函数是异步的,里面可能有等待操作」。
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:.2f}秒")
if __name__ == '__main__':
asyncio.run(main())
运行结果:
开始下载图片:风景.jpg
开始下载图片:人物.jpg
开始下载图片:宠物.jpg
图片下载完成
图片下载完成
图片下载完成
总耗时:2.00秒
总耗时从 6 秒降到了 2 秒——这就是异步的魔力:三个任务的「等待时间」重叠了。
asyncio.gather() 接收多个协程,并发执行它们,等所有协程都完成后返回结果列表。核心机制如下图:
asyncio.gather(task1, task2, task3)
│
▼
┌─────────┐
│ 事件循环 │ ← 单线程内的"调度助手"
└────┬────┘
│ 同时发起
┌────┼────┬────┐
▼ ▼ ▼ ▼
task1 task2 task3 ← 三个IO任务并发等待
│ │ │
└────┼────┘
▼
全部完成,返回结果 ← 总耗时 ≈ 最慢的那个任务
1.3 同步 vs 异步:什么时候该用哪个?
很多人觉得「异步比同步好,所有场景都该用异步」——这其实是误区。同步和异步各有适用场景,选错了反而会适得其反。
| 对比维度 | 同步 | 异步 |
|---|---|---|
| 执行方式 | 串行,一个做完再做下一个 | 并发,等待时切换做其他 |
| 适用场景 | CPU 密集型任务 | IO 密集型任务 |
| 典型场景 | 大量计算、本地文件读写 | 网络请求、数据库查询、文件下载 |
| 学习成本 | 低,符合直觉 | 高,需理解事件循环、协程 |
| 并发能力 | 靠多线程/多进程 | 单线程即可 |
| 典型库 | requests、sqlite3 |
aiohttp、asyncpg |
举个例子:
- 如果你写一个「计算 1 到 100000 的质数」的程序(CPU 密集),用异步没用——因为 CPU 一直在工作,没有等待间隙,异步无法切换任务
- 如果你写一个「爬取 100 个网页」的程序(IO 密集),用异步就很合适——因为爬网页时大部分时间在等服务器响应,异步可以在等待时爬其他网页
1.4 异步常见误区
误区一:异步 ≠ 多线程/多进程
异步是「单线程下的任务调度」,而多线程/多进程是「真正的并行执行」。异步不能解决 CPU 密集型任务的效率问题,因为它本质还是单线程。
误区二:不是所有库都支持异步
很多 Python 库是同步的(比如 requests、sqlite3),在异步函数里用这些库会阻塞线程——必须用对应的异步库:
| 同步库 | 异步替代 | 用途 |
|---|---|---|
requests |
aiohttp |
HTTP 请求 |
sqlite3 |
aiosqlite |
SQLite 数据库 |
psycopg2 |
asyncpg |
PostgreSQL |
redis |
redis.asyncio |
Redis 缓存 |
误区三:异步代码不一定比同步快
如果任务量很小,或者 IO 等待时间极短(比如本地文件读写),异步的「调度开销」可能比节省的等待时间还多,反而更慢。
一句话总结:异步就是不要让 CPU 闲下来——当任务需要大量等待(IO 密集),用异步(
async/await+asyncio);当任务需要大量计算(CPU 密集),用同步 + 多线程/多进程。
1.5 异步封装:让同步代码跑异步任务
在实际项目中,agent.invoke() 本身是同步调用方法,会阻塞线程。如果你已经有了异步的 Agent(agent.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
注意:
asyncio.run()每次调用都会创建新的事件循环。如果需要在同一程序中多次调用异步代码,建议复用同一个事件循环,而非反复创建销毁。
Block 2:MCP 协议——智能体世界的「USB 接口」
2.1 什么是 MCP?为什么需要它?
先想一个问题:FastAPI 中有很多接口——有些是自己写的,有些是互联网上的(比如快递查询、天气预报),接口的相互调用让项目功能更丰富。在智能体的世界中,也有同样的需求,我们管它叫 MCP(Model Context Protocol,模型上下文协议)。
MCP 用于标准化模型与外部工具、数据源和上下文环境之间的交互方式。它的目标是:让不同系统之间共享统一的模型上下文接口,使智能体能在不同的运行环境中调用工具时,无须关心底层传输细节。
通俗地说,MCP 就像智能体世界的「USB 协议」——无论工具在本地还是异地,使用 Python、Node.js 或 HTTP,都可以通过统一接口被智能体调用。
没有 MCP 的世界 有 MCP 的世界
┌──────────┐ ┌──────────┐
│ Agent │ │ Agent │
└────┬─────┘ └────┬─────┘
│ 各写各的胶水代码 │ 统一 MCP 接口
├─→ HTTP API (自定义格式) │
├─→ Shell 脚本 (自定义调用) ┌────┴─────┐
├─→ gRPC (自定义proto) │ MCP │
└─→ WebSocket (自定义协议) │ Protocol │
┌────┴────┬────┴────┐
▼ ▼ ▼
工具A 工具B 工具C
(stdio) (http) (python)
MCP 诞生时间线:2024 年 11 月 25 日,Anthropic 正式发布并开源 MCP 协议。
2.2 MCP 的核心架构:Client-Server
MCP 本质上是一个客户端-服务端协议(Client-Server 协议):
- MCP 服务端(Server):负责定义并向网络暴露可用的工具
- MCP 客户端(Client):运行在智能体端,负责发现、加载并调用这些远程工具
LangChain 官方提供了 langchain-mcp-adapters 适配库来支持 MCP,它可以将 MCP 工具无缝集成到 LangChain Agent 中。
2.3 实战:从零搭建 MCP 工具链
第一步:安装依赖
pip install mcp fastmcp langchain-mcp-adapters
第二步:创建数学工具 MCP 服务端(stdio 传输)
创建文件 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()
def multiply(a: int, b: int) -> int:
"""Multiply two numbers"""
return a * b
if __name__ == "__main__":
# transport="stdio" 表示通过标准输入输出进行通信
# 客户端会作为子进程拉起此服务端
mcp.run(transport="stdio")
关键点解析:
FastMCP("Math"):创建一个名为 “Math” 的 MCP 服务实例@mcp.tool():将普通函数注册为 MCP 工具,函数名即工具名,docstring 即工具描述transport="stdio":使用标准输入输出通信,适合本地子进程模式——客户端会自动拉起这个服务端作为子进程
第三步:创建天气工具 MCP 服务端(HTTP 传输)
创建文件 weather_mcp_server.py,此服务端调用免费的 Open-Meteo API(无需 API Key),并采用 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"
# 工具一:将城市名解析为经纬度
@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"),
}
# 工具二:根据经纬度查询当前天气
@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"),
}
# 工具三:组合工具——城市名 → 天气(内部先地理编码再查询天气)
@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"])
# **g 解包字典:等价于手动传 name=g["name"], lat=g["lat"] ...
return {**g, **w}
if __name__ == "__main__":
print("Starting Weather MCP Server (streamable-http) on http://localhost:8000/mcp ...")
# transport="streamable-http" 表示使用 HTTP 协议暴露服务
# 端点默认路径是 /mcp,端口通常为 8000
mcp.run(transport="streamable-http")
关于 **g 字典解包:上面 get_current_weather_by_city 中用到了 {**g, **w},这是 Python 的字典解包语法。原理如下:
g = {"name": "张三", "age": 20}
def func(name, age):
print(name, age)
# 等价写法1:手动传关键字
func(name="张三", age=20)
# 等价写法2:**g 自动拆解字典,等价于 func(name="张三", age=20)
func(**g)
# 同理,{**g, **w} 就是把两个字典合并成一个新字典
启动天气服务端:
python weather_mcp_server.py
# 输出: Starting Weather MCP Server (streamable-http) on http://localhost:8000/mcp ...
第四步:创建 MCP 客户端并集成智能体
创建文件 client_mcp_demo.py,该客户端将同时连接上述两个 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()
# 连接 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)。"
)
# 任务1:数学计算
result1 = asyncio.run(agent.ainvoke({
"messages": [
{
"role": "user",
"content": "请帮我计算 (3 + 5) × 12 的结果"
}
]
}))
for message in result1['messages']:
print(type(message).__name__)
print(message.content)
# 任务2:天气查询
result2 = asyncio.run(agent.ainvoke({
"messages": [
HumanMessage(content="今天郑州天气如何?")
]
}))
for message in result2['messages']:
print(type(message).__name__)
print(message.content)
运行输出:
已加载的工具: ['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
2.4 整体架构与执行流程梳理
上述示例构建了一个完整的 MCP 集成环境,其架构与数据流如下:
┌──────────────────────────────────────────────────────────────┐
│ 客户端(Agent 侧) │
│ │
│ ┌──────────────┐ ┌──────────────────────┐ │
│ │ MultiServer │ │ create_agent │ │
│ │ MCPClient │───→│ (model + tools) │ │
│ └──────┬───────┘ └──────────────────────┘ │
│ │ │
│ ┌────┴────┐ │
│ ▼ ▼ │
│ stdio streamable_http │
│ │ │ │
└────┼─────────┼───────────────────────────────────────────────┘
│ │
▼ ▼
┌─────────┐ ┌──────────────┐
│Math │ │Weather │
│Server │ │Server │
│(stdio) │ │(HTTP :8000) │
│ │ │ │
│add() │ │geocode_city()│
│multiply()│ │get_weather() │
└─────────┘ └──────────────┘
三个核心组件详解:
| 组件 | 文件 | 传输方式 | 暴露工具 | 启动方式 |
|---|---|---|---|---|
| 数学服务端 | math_mcp_server.py |
stdio | add、multiply |
客户端自动拉起子进程 |
| 天气服务端 | weather_mcp_server.py |
streamable-http | geocode_city、get_current_weather、get_current_weather_by_city |
手动 python 启动 |
| 客户端+Agent | client_mcp_demo.py |
— | 聚合两个服务端的工具 | python 运行 |
执行流程:
- 客户端通过
MultiServerMCPClient同时连接两个服务端 - 调用
client.get_tools()拉取所有工具定义 - 将工具交给
create_agent注册 - 智能体运行时采用 ReAct 模式,根据用户问题自动选择工具并调用
- Math 工具通过 stdio 本地通信,Weather 工具通过 HTTP 远程调用
2.5 两种传输方式对比
| 对比维度 | stdio | streamable-http |
|---|---|---|
| 通信方式 | 标准输入输出 | HTTP 请求 |
| 部署形态 | 本地子进程 | 网络服务 |
| 跨机器 | 不支持 | 支持 |
| 启动方式 | 客户端自动拉起 | 手动启动服务 |
| 适用场景 | 本地开发、轻量工具 | 生产部署、远程调用 |
| 性能 | 高(进程间直接通信) | 中(HTTP 开销) |
Block 3:预置中间件——开箱即用的生产保障
3.1 为什么需要中间件?
从 Block 2 我们知道,MCP 让 Agent 能调用外部工具。但生产环境中,Agent 还面临很多问题:
- 对话太长,Token 超限怎么办?→ 需要自动摘要
- Agent 要执行删库操作,谁审批?→ 需要人工干预
- 模型调用失败,要不要重试?→ 需要重试机制
- 工具调用太频繁,成本怎么控?→ 需要限流
这些需求有一个共同特点:它们与 Agent 的核心业务逻辑无关,而是在执行流程的特定阶段自动触发的横切关注点。这就是中间件的用武之地。
中间件是什么? 中间件是在智能体生命周期中自动触发的钩子函数,帮我们进行权限校验、日志记录、消息压缩、人工介入等操作——跟业务逻辑解耦,插拔即用。
3.2 LangChain 预置中间件全景
LangChain 官方提供了一系列开箱即用的预置中间件,从可靠性、可观测性、成本控制与安全治理等多个维度,为智能体系统提供保障:
| 中间件 | 核心能力 | 典型场景 |
|---|---|---|
| Summarization | 对话历史自动压缩 | 长对话 Token 超限 |
| 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.3 示例一:SummarizationMiddleware(对话摘要)
用途:自动对对话历史进行压缩,避免超出模型的上下文窗口限制,同时尽力维持对话的连贯性。
关键机制:在模型调用前,检查上下文 Token 数量,若超过预设阈值则触发摘要进程,将早期消息压缩为简洁的摘要,并可配置保留最近 N 条原始消息以维持细节。
from langchain.agents import create_agent
from langchain.agents.middleware import SummarizationMiddleware
from langchain_core.messages import HumanMessage
from utils.model_factory import get_deepseek_model
llm = get_deepseek_model(temperature=0.7)
summarization_middleware = SummarizationMiddleware(
model=llm,
trigger=('tokens', 1000), # 达到 1000 token 触发摘要
keep=('messages', 2), # 保留最近 2 条消息
summary_prompt="用20个字以内概括要点。", # 自定义摘要提示词(不设则用默认)
)
agent = create_agent(
model=llm,
system_prompt="请使用30字以内回答问题",
middleware=[summarization_middleware],
tools=[]
)
conversation = [
"RAG是什么?",
"有何用途?",
"主要缺点?",
"一句话总结"
]
state = {"messages": []}
for i, question in enumerate(conversation, start=1):
human_message = HumanMessage(content=f"{question}")
state = agent.invoke({
"messages": state['messages'] + [human_message]
})
answer = state['messages'][-1].content
print(f"\n第 {i} 轮")
print(f"Q: {question}")
print(f"A: {answer}")
运行结果:
第 1 轮
Q: RAG是什么?
A: RAG是检索增强生成技术。
第 2 轮
Q: 有何用途?
A: 提升回答准确性和知识广度。
第 3 轮
Q: 主要缺点?
A: 依赖检索质量,可能增加延迟。
第 4 轮
Q: 一句话总结
A: RAG通过检索增强生成,提升准确但依赖检索质量。
对话结束。
SummarizationMiddleware 关键参数说明:
| 参数 | 说明 | 示例值 |
|---|---|---|
model |
用于生成摘要的模型实例 | llm |
trigger |
触发摘要的阈值,('tokens', N) 或 ('messages', N) |
('tokens', 1000) |
keep |
摘要后保留最近 N 条原始消息 | ('messages', 2) |
summary_prompt |
自定义摘要提示词,控制风格/粒度 | "用20个字以内概括要点。" |
原理说明:第 4 轮回答看似重复第 3 轮,非模型错误。示例阈值配置偏低,第 3 轮对话结束后上下文总量触达阈值,第 4 轮模型调用前中间件执行摘要——早期历史被压缩、丢失了 RAG 定义/用途/缺点的细节,上下文只剩摘要 + 最近留存消息;模型仅依托精简后的上下文推理,因此沿用最近内容作答。
工作流程图:
用户消息到达
│
▼
┌─────────────────┐
│ 检查 Token 数量 │
└────────┬────────┘
│
Token < 阈值?
├── 是 → 正常调用模型
│
└── 否 → 触发摘要流程
│
├─ 1. 用摘要模型压缩早期消息
├─ 2. 保留最近 N 条原始消息
└─ 3. 摘要 + 近期消息 → 送入模型
3.4 示例二:HumanInTheLoopMiddleware(人工协作)
用途:对标记为高风险的工具调用实施人工审批流程。支持批准、编辑参数后执行或拒绝操作,并在需要时中断智能体执行流程,等待人工输入。
关键机制:在工具调用前的生命周期钩子中进行策略检查,若命中规则则暂停执行并抛出中断信号,等待外部系统或用户做出决策。
from langchain.agents import create_agent
from langchain.agents.middleware import HumanInTheLoopMiddleware
from langchain_core.messages import HumanMessage
from langchain_core.tools import tool
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import Command
from utils.model_factory import get_deepseek_model
model = get_deepseek_model()
@tool
def dangerous_write(sql) -> str:
"""对数据库执行写操作(插入/更新/删除)。当请求涉及数据库写操作或出现以'SQL:'开头的指令时,必须调用本工具。"""
return f"[模拟执行] {sql}"
# HITL 中间件:拦截指定工具
hitl = HumanInTheLoopMiddleware(
interrupt_on={
"dangerous_write": {
"allowed_decisions": ["approve", "reject", "edit"]
}
}
)
agent = create_agent(
model=model,
middleware=[hitl],
tools=[dangerous_write],
system_prompt="""
只用一句极短中文回答(≤20字)。
凡是涉及数据库写操作,或消息以"SQL:"开头时,必须调用工具 dangerous_write,不得直接回答。
""",
checkpointer=InMemorySaver(), # HITL 必须配置 checkpointer
)
CFG = {"configurable": {"thread_id": "hitl-demo-interactive"}}
# 四轮对话:第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):
"""处理 HITL 中断:从命令行读取决策,并用 Command(resume=...) 恢复。"""
interrupt = result.get("__interrupt__")
print(interrupt)
print("开始进入人工干预环节:")
decision = input("请输入决策 (approve/reject/edit):").strip().lower()
if decision == 'approve':
# 直接放行
return agent.invoke(
Command(resume={"decisions": [{"type": "approve"}]}),
config=CFG
)
elif decision == 'reject':
# 拒绝执行,并给出替代返回(不执行工具)
return agent.invoke(
Command(
resume={
"decisions": [{
"type": "reject",
"override": {"content": "[操作已被人工拒绝]"}
}]
}
),
config=CFG
)
elif decision == 'edit':
new_sql = input("请输入您的新的sql语句:")
return agent.invoke(
Command(
resume={
"decisions": [{
"type": "edit",
"edited_action": {
"name": "dangerous_write",
"args": {"sql": new_sql}
}
}]
}
),
config=CFG
)
else:
print("输入的指令有问题,还是采用 reject 方案")
return agent.invoke(
Command(
resume={
"decisions": [{
"type": "reject",
"override": {"content": "[操作已被人工拒绝]"}
}]
}
),
config=CFG
)
for i, question in enumerate(conversation, start=1):
print(f"\n第 {i} 轮")
print(f"Q: {question}")
human_message = HumanMessage(content=f"{question}")
state = agent.invoke({
"messages": [human_message]
}, config=CFG)
# 命中中断:让你做决策 → 再 resume
if state.get("__interrupt__"):
state = handle_interrupt(state)
answer = state['messages'][-1].content
print(f"A: {answer}")
HITL 工作流程图:
用户提问 → Agent 推理 → 需要调用 dangerous_write?
│
是 ──┤── 否 → 直接返回
│
┌─────────▼──────────┐
│ HITL 中断执行 │
│ 等待人工决策 │
└─────────┬──────────┘
│
┌───────────────┼───────────────┐
▼ ▼ ▼
approve edit reject
│ │ │
▼ ▼ ▼
执行原SQL 执行修改后SQL 返回拒绝信息
│ │ │
└───────────────┼───────────────┘
▼
Agent 继续执行
HITL 常见错误及解决方案
错误现象:运行时报如下错误:
openai.BadRequestError: Error code: 400 - {
'error': {
'message': "An assistant message with 'tool_calls' must be followed by
tool messages responding to each 'tool_call_id'. (insufficient tool messages
following tool_calls message)",
'type': 'invalid_request_error'
}
}
错误原因:在使用了 checkpointer(如 InMemorySaver())的 Agent 中,checkpointer 本身就具有历史消息保存功能。如果在代码中手动拼接历史消息:
# 错误写法——手动拼接历史消息
state = agent.invoke({
"messages": state['messages'] + [human_message]
}, config=CFG)
会导致消息重复或格式错乱,因为 checkpointer 已经在内部维护了完整的历史记录。
解决方案:每次只传入新消息,让 checkpointer 自动管理历史:
# 正确写法——只传入新消息
state = agent.invoke({
"messages": [human_message]
}, config=CFG)
总结:LangChain 1.0 提供的预置中间件生态丰富而强大,覆盖了日志追踪、对话管理、安全审查、人工干预等关键生产需求。它们在智能体运行时扮演着不同的角色:一部分专注于性能优化与资源管理,另一部分则侧重于安全、合规与流程控制。本节仅通过对话摘要与人工协作两个典型案例展示了预置中间件的使用方法,其他中间件的使用方法类似,可自行探索。
Block 4:自定义中间件——从装饰器到类的进阶之路
4.1 两种自定义方式
预置中间件覆盖了通用场景,但生产环境中总有定制需求。LangChain 提供了两种自定义中间件的方式:
| 方式 | 适用场景 | 复杂度 |
|---|---|---|
| 基于装饰器(Decorator-based) | 逻辑简单、只需在一个钩子上注入功能 | 低 |
| 基于类(Class-based) | 功能复杂、需要在多个钩子间维护状态或协调逻辑 | 中 |
4.2 装饰器模式:三类钩子
在装饰器模式下,LangChain 官方将可用的钩子分为三类,以适应不同的拦截与处理需求。
第一类:节点式(Node-style)钩子
在特定的执行点运行,按时间顺序触发(before → after),不会改变函数调用结构:
| 钩子 | 触发时机 | 典型场景 |
|---|---|---|
@before_agent |
智能体执行前 | 初始化、上下文检查、权限验证 |
@before_model |
每次模型调用前 | 修改输入、记录日志、过滤 |
@after_model |
每次模型响应后 | 分析输出、检测风险词、结果修正 |
@after_agent |
智能体执行完成后 | 收尾工作、日志上报、状态持久化 |
第二类:包裹式(Wrap-style)钩子
用于包裹实际执行过程,在调用前后插入逻辑,可直接控制函数执行(甚至中断或替换结果):
| 钩子 | 触发时机 | 典型场景 |
|---|---|---|
@wrap_model_call |
每次模型调用时 | 计时、限流、重试、异常兜底 |
@wrap_tool_call |
每次工具调用时 | 审计、权限校验、记录调用轨迹 |
第三类:便捷装饰器(Convenience decorator)
常用于根据上下文动态调整系统提示:
| 钩子 | 触发时机 | 典型场景 |
|---|---|---|
@dynamic_prompt |
模型调用前 | 千人千面 Prompt、个性化身份/语言/风格 |
@dynamic_prompt相当于对@wrap_model_call的简化封装,专用于在模型调用前动态修改 Prompt。
4.3 装饰器实战:完整示例
下面通过一个完整示例展示所有三类装饰器中间件的协同使用:
import time
from typing import Callable
from langchain.agents import create_agent, AgentState
from langchain.agents.middleware import (
before_model, after_model, before_agent, after_agent,
wrap_model_call, dynamic_prompt,
ModelRequest, ModelResponse
)
from langchain_core.messages import HumanMessage
from langgraph.runtime import Runtime
from utils.model_factory import get_deepseek_model
llm = get_deepseek_model()
# ========= 节点式钩子 =========
@before_model # 在 model 调用之前触发
def log_before_model(state: AgentState, runtime: Runtime):
print("[before_model] 准备进行模型调用,当前消息数:", len(state['messages']))
return None
@after_model # 在 model 调用之后触发
def log_after_model(state: AgentState, runtime: Runtime):
print("[after_model] 模型调用完成,当前消息数:", len(state['messages']))
return None
@before_agent # 在 agent 调用之前触发
def log_before_agent(state: AgentState, runtime: Runtime):
print("[before_agent] 智能体开始执行,当前消息数:", len(state['messages']))
return None
@after_agent # 在 agent 调用之后触发
def log_after_agent(state: AgentState, runtime: Runtime):
print("[after_agent] 智能体执行完成,当前消息数:", len(state['messages']))
return None
# ========= 包裹式钩子:为模型调用加重试与耗时统计 =========
@wrap_model_call
def retry_and_timing(request: ModelRequest,
handler: Callable[[ModelRequest], ModelResponse]):
start = time.time()
try:
return handler(request)
except Exception as e:
print(f"[wrap_model_call] 模型调用失败: {e}")
return None
finally:
end = time.time()
cost = (end - start) * 1000
print(f"[wrap_model_call] 本次模型调用耗时:{cost:.0f} ms")
# ========= 便捷装饰器:动态系统提示 =========
# 在每次 LLM 调用前动态生成系统提示词,支持读取运行时上下文(runtime.context)
# 实现千人千面的系统 Prompt
@dynamic_prompt
def personalized_prompt(req: ModelRequest) -> str:
runtime = req.runtime
user_id = "访客"
if runtime and getattr(runtime, "context", None):
user_id = runtime.context.get("user_id", "访客")
return (f"你是一名贴心的中文助手,正在为用户{user_id}提供帮助。"
"回答时要简洁、自然,回答问题的时候先说用户的user_id,再回答问题。")
agent = create_agent(
model=llm,
tools=[],
middleware=[
log_before_model,
log_after_model,
log_before_agent,
log_after_agent,
retry_and_timing,
personalized_prompt
],
system_prompt="回答问题要简洁,最好不要超过30个字"
)
result = agent.invoke(
{"messages": [HumanMessage("用一句话解释 LangGraph 是什么。")]},
context={"user_id": "alice"},
)
for message in result['messages']:
print(type(message).__name__)
print(message.content)
print("~" * 30)
执行顺序详解:
当 Agent 处理一条用户消息时,各钩子的触发顺序如下:
agent.invoke()
│
▼
@before_agent ──────→ 智能体开始执行(日志/初始化)
│
▼
┌──────────────────────────────────────────┐
│ 模型调用阶段 │
│ │
│ @before_model ──→ 模型调用前(日志/过滤) │
│ │ │
│ ▼ │
│ @dynamic_prompt ──→ 动态生成 Prompt │
│ │ │
│ ▼ │
│ @wrap_model_call ──→ 包裹调用(计时/重试)│
│ │ │
│ ├──→ handler(request) 实际调用模型 │
│ │ │
│ ▼ │
│ @after_model ──→ 模型调用后(结果分析) │
└──────────────────────────────────────────┘
│
▼
@after_agent ────────→ 智能体执行完成(收尾/持久化)
│
▼
返回结果
关于 state 参数:在钩子函数中看到
state参数,它就是历史数据的意思——包含了到目前为止的所有消息和上下文信息。
4.4 关于「第二次回答有 3 条消息」的解释
在使用中间件时,你可能会发现一个现象:第二次调用 Agent 后,返回的 messages 列表有 3 条消息,而不是预期的 1 条。原因如下:
第一次 invoke:
messages = [HumanMessage] → 1条
第二次 invoke(state 是隔离的,上次的消息被清空):
messages = [HumanMessage] → 第1条:用户新提问
messages += [AIMessage] → 第2条:模型生成的回答
messages += [AIMessage (修改后)] → 第3条:被 @after_model 拦截后修改的回答
注意:消息不会覆盖,只会追加,所以有 3 条
原因:两次 agent.invoke 之间的 state 是隔离的,之前的回答被清空了。中间件(如 @after_model)修改了 AIMessage 后,由于消息不会覆盖只会追加,所以列表中会多出一条修改后的消息。
4.5 基于类的中间件
当中间件逻辑变得复杂,需要组合使用多个钩子,或需要维护内部状态以支持可配置性与复用时,基于类(Class-based)的中间件实现方式便成为更合适的选择。
其核心钩子类型与基于装饰器的钩子类型类似:
- 节点式(Node-style)钩子:在执行的固定阶段顺序触发,适用于日志记录、状态校验、上下文更新等场景
- 包裹式(Wrap-style)钩子:用于包裹实际的调用过程,能够深度控制执行流程,典型应用包括实现重试、性能统计、结果缓存等场景
import time
from typing import Callable
from langchain.agents import AgentState, create_agent
from langchain.agents.middleware import AgentMiddleware, ModelRequest, ModelResponse
from langchain_core.messages import HumanMessage
from langgraph.runtime import Runtime
from utils.model_factory import get_deepseek_model
class PolicyGuardMiddleware(AgentMiddleware):
"""在关键节点打印日志"""
def before_model(self, state: AgentState, runtime: Runtime):
print(f"[before_model] 准备进行模型调用,当前消息数:{len(state['messages'])}")
return None
def after_model(self, state: AgentState, runtime: Runtime):
print(f"[after_model] 模型调用完成,当前消息数:{len(state['messages'])}")
return None
def before_agent(self, state: AgentState, runtime: Runtime):
print(f"[before_agent] 智能体开始执行,当前消息数:{len(state['messages'])}")
return None
def after_agent(self, state: AgentState, runtime: Runtime):
print(f"[after_agent] 智能体执行完成,当前消息数:{len(state['messages'])}")
return None
class RetryAndMetricsMiddleware(AgentMiddleware):
"""包裹每次模型调用,做退避重试与耗时统计"""
def wrap_model_call(self, request: ModelRequest,
handler: Callable[[ModelRequest], ModelResponse]):
start = time.time()
try:
return handler(request)
except Exception as e:
print(f"[wrap_model_call] 模型调用失败: {e}")
finally:
end = time.time()
cost = (end - start) * 1000
print(f"[wrap_model_call] 本次模型调用耗时:{cost:.0f} ms")
agent = create_agent(
model=get_deepseek_model(),
middleware=[
PolicyGuardMiddleware(), # 注意:要 new 对象,不是传类名
RetryAndMetricsMiddleware()
],
tools=[],
system_prompt="回答问题要简洁,不要超过40个字"
)
result = agent.invoke(
{"messages": [HumanMessage("用一句话解释 LangGraph 是什么。")]},
context={"user_id": "alice"},
)
for message in result['messages']:
print(type(message).__name__)
print(message.content)
print("~" * 30)
类中间件注意事项:
- 在
middleware列表中传参时,要 实例化对象(PolicyGuardMiddleware()),而不是传类名 - 类名不重要,重要的是必须继承自
AgentMiddleware - 类中的方法名不能随意起,必须使用固定的名字(如
before_agent、after_model等) - 类中的方法必须传递
self参数
4.6 多中间件执行顺序:洋葱模型
当多个中间件同时存在时,其执行顺序遵循以下规则——这就是经典的洋葱模型:
请求进入
│
▼
┌─── before_agent (中间件A) ───┐
│ │
│ ┌── before_agent (中间件B) │
│ │ │
│ │ ┌── before_model (A) │
│ │ │ │
│ │ │ ┌── 模型调用 ──┐ │
│ │ │ │ │ │
│ │ │ │ LLM 实际执行 │ │
│ │ │ │ │ │
│ │ │ └───────────────┘ │
│ │ │ │
│ │ └── after_model (A) ──┘
│ │ │
│ └── after_agent (中间件B)──┘
│ │
└─── after_agent (中间件A) ────┘
│
▼
请求返回
执行顺序规则:
| 钩子类型 | 执行顺序 | 说明 |
|---|---|---|
before_* 钩子 |
按注册顺序(从左到右) | 先注册的先执行 |
| 包裹式钩子 | 外层先于内层 | 洋葱模型,层层包裹 |
after_* 钩子 |
按注册逆序(从右到左) | 先注册的后执行 |
这三类中间件并非互斥的,而是可以协同使用的。例如,可以先用
@before_model进行输入校验,再通过@wrap_model_call实现调用重试,从而构建出功能完善且职责清晰的中间件链。
Block 5:调优答疑与总结
5.1 FAQ 常见问题
Q1:异步代码中能否调用同步库(如 requests)?
不能直接调用。同步库会阻塞事件循环,导致整个异步程序卡住。解决方案:
- 换用异步替代库(
requests→aiohttp) - 或用
asyncio.to_thread()将同步调用包装到线程中执行
Q2:MCP 的 stdio 和 streamable-http 该选哪个?
- 本地开发/轻量工具:用 stdio,客户端自动拉起子进程,零配置
- 生产部署/远程调用:用 streamable-http,支持跨机器访问、多客户端共享
Q3:SummarizationMiddleware 的阈值怎么设?
取决于模型的上下文窗口大小。建议设为模型上下文窗口的 60%-70%,留出空间给新消息。例如模型支持 8K Token,阈值设为 5000-6000。
Q4:HITL 中断后如何恢复执行?
通过 Command(resume={"decisions": [...]}) 恢复。三种决策类型:
approve:批准并继续执行reject:拒绝执行,返回替代信息edit:修改工具调用的输入参数后继续执行
Q5:中间件和工具有什么区别?
| 对比维度 | 中间件 | 工具 |
|---|---|---|
| 触发方式 | 生命周期自动触发 | Agent 主动调用 |
| 与业务关系 | 无关(横切关注点) | 直接相关 |
| 典型场景 | 日志、安全、摘要 | 查天气、算数学 |
| 返回值 | 修改状态/拦截流程 | 返回工具执行结果 |
Q6:装饰器中间件和类中间件怎么选?
- 逻辑简单、单一钩子:用装饰器,代码简洁
- 逻辑复杂、多钩子组合、需维护状态:用类,结构清晰、可复用
Q7:多个中间件的执行顺序如何控制?
通过 middleware 列表的注册顺序控制。before_* 按列表顺序执行,after_* 按列表逆序执行,包裹式遵循洋葱模型。
5.2 中间件分类总结
按实现方式分三类:
- 预置中间件:LangChain 官方提供的开箱即用中间件,如 SummarizationMiddleware、HumanInTheLoopMiddleware 等
- 装饰器中间件:用一个函数 + 装饰器(如
@before_model、@after_model)实现,适合简单场景 - 基于类的中间件:继承
AgentMiddleware,在类中编写对应的触发方法,适合复杂场景
按钩子类型分三类:
- 节点式钩子:关注「何时触发」,在特定执行点轻量地插入逻辑,不改变原有调用流程
- 包裹式钩子:关注「如何执行」,能够包裹并干预完整的调用过程,具备中断、修改等强大控制能力
- 便捷装饰器:关注「效率」,为动态修改 Prompt 等常见场景提供快速实现方案
5.3 核心 API 速查表
| API / 装饰器 | 所属 | 作用 |
|---|---|---|
asyncio.run(coro) |
Python 标准库 | 创建事件循环并运行协程 |
asyncio.gather(*coros) |
Python 标准库 | 并发执行多个协程 |
FastMCP(name) |
mcp 库 | 创建 MCP 服务端实例 |
@mcp.tool() |
mcp 库 | 注册函数为 MCP 工具 |
mcp.run(transport=...) |
mcp 库 | 启动 MCP 服务(stdio / streamable-http) |
MultiServerMCPClient(config) |
langchain-mcp-adapters | 连接多个 MCP 服务端 |
client.get_tools() |
langchain-mcp-adapters | 拉取所有服务端的工具 |
SummarizationMiddleware(...) |
LangChain 预置 | 对话历史自动摘要 |
HumanInTheLoopMiddleware(...) |
LangChain 预置 | 敏感操作人工审批 |
@before_agent |
装饰器钩子 | Agent 执行前触发 |
@before_model |
装饰器钩子 | 模型调用前触发 |
@after_model |
装饰器钩子 | 模型响应后触发 |
@after_agent |
装饰器钩子 | Agent 执行后触发 |
@wrap_model_call |
包裹式钩子 | 包裹模型调用(计时/重试/限流) |
@wrap_tool_call |
包裹式钩子 | 包裹工具调用(审计/权限校验) |
@dynamic_prompt |
便捷装饰器 | 动态生成系统 Prompt |
AgentMiddleware |
基类 | 类中间件基类 |
Command(resume=...) |
langgraph.types | 恢复 HITL 中断后的执行 |
InMemorySaver() |
langgraph.checkpoint | 内存级状态持久化(HITL 必需) |
5.4 落地场景
| 场景 | 推荐方案 | 核心中间件 |
|---|---|---|
| 高并发 API 服务 | 异步 Agent + MCP 远程工具 | @wrap_model_call(限流+重试) |
| 长对话客服机器人 | 异步 + Summarization | SummarizationMiddleware |
| 数据库运维 Agent | MCP + HITL | HumanInTheLoopMiddleware |
| 多租户个性化助手 | @dynamic_prompt |
千人千面 Prompt |
| 企业级 Agent 平台 | 全套中间件 + 类中间件 | 日志+安全+摘要+限流 |
5.5 本章核心知识点
- 异步编程:
async/await+asyncio让 IO 密集型任务在单线程内并发执行,asyncio.gather()是并发执行多个协程的核心 API - MCP 协议:智能体世界的「USB 接口」,Client-Server 架构,支持 stdio 和 streamable-http 两种传输方式,
MultiServerMCPClient可同时连接多个服务端 - 预置中间件:LangChain 提供 12 种开箱即用的中间件,覆盖摘要、人工审批、限流、重试等生产需求
- 自定义中间件:装饰器模式适合简单场景(三类钩子:节点式/包裹式/便捷装饰器),类模式适合复杂场景(继承
AgentMiddleware) - 执行顺序:
before_*按注册顺序执行,after_*按逆序执行,包裹式遵循洋葱模型
下一篇预告:我们将进入 LangChain 系列的压轴篇——从零搭建一个完整的企业级 AI Agent 项目,整合 RAG、Tool、MCP、中间件等所有能力,形成端到端的实战案例。敬请期待。
本文所有代码均经过验证,可直接复制运行。如有问题欢迎评论区交流。
更多推荐


所有评论(0)