python-sdk 客户端订阅完全指南:用 client.listen() 实时监听 MCP 资源与目录变化
【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址: https://gitcode.com/gh_mirrors/pythonsd/python-sdk
导读
MCP(Model Context Protocol)服务器的目录(catalog)并非一成不变:工具可能在运行时才出现,某个资源 URI 背后的内容也会不断变化。客户端如何第一时间感知这些变化?答案是client.listen(...)——一次subscriptions/listen请求,而该请求的响应本身就是一条长期打开的流(stream),源源不断地推送客户端主动订阅的变更通知。本文基于 python-sdk 的官方文档与源码,系统讲解客户端这一端的完整故事:如何打开订阅流、如何在主流程之外并行监听、如何处理流的各种结束方式,以及 SDK 底层的去重、过滤与多路复用实现,让你能写出可靠、不阻塞、可优雅恢复的订阅监听代码。
订阅的本质:一次请求,一条流
在 src/mcp/client/subscriptions.py 的模块注释中,SDK 给出了清晰的定义:
listen()opens the stream as an async context manager: entering waits for the server's acknowledgment, iteration yields typed change events, a graceful server close ends the loop, and an abrupt drop raisesSubscriptionLost. There is no replay and no automatic re-listen.
核心要点有三个:
- 进入即发送请求并等待确认:
async with client.listen(...)进入上下文管理器时,SDK 会发送subscriptions/listen请求,并把你的关键字参数构造成订阅过滤器(SubscriptionFilter),然后阻塞等待服务器的确认(acknowledgment)。因此,当代码块真正开始执行时,这条流已经是"活的"了——之后发布的所有变更都会送达。 - 迭代即消费事件:对订阅对象做
async for event in sub迭代,会逐个收到类型化(typed)的事件对象。 - 没有重放,也没有自动重连:流一旦结束就彻底结束,客户端需要自行重新 listen 并重新获取数据。
从协议版本上看,subscriptions/listen是 2026-07-28 版本(SEP-2575)引入的能力。如果连接协商出的协议版本早于该版本,SDK 会直接抛出类型化的ListenNotSupportedError,引导你改用subscribe_resource()等旧版机制(详见下文"进入可能抛出的异常"一节)。
监听一条流:四个事件类型与订阅句柄
本文的所有示例都围绕文档教程中构建的 sprint-board(冲刺看板)服务器展开。首先来看最核心的监听代码,它来自 docs_src/subscriptions/tutorial003.py:
from mcp import Client from mcp.client.subscriptions import ResourceUpdated, ToolsListChanged from mcp.types import TextResourceContents BOARD = "board://sprint" async def read_board(client: Client, uri: str = BOARD) -> str: [contents] = (await client.read_resource(uri)).contents assert isinstance(contents, TextResourceContents) return contents.text async def follow_board(client: Client) -> None: async with client.listen(tools_list_changed=True, resource_subscriptions=[BOARD]) as sub: async for event in sub: match event: case ResourceUpdated(uri=uri): print(await read_board(client, uri)) case ToolsListChanged(): tools = await client.list_tools() print("tools:", [tool.name for tool in tools.tools]) case _: pass # kinds the filter did not ask for never arrive async def main() -> None: async with Client("http://localhost:8000/mcp") as client: await follow_board(client)过滤器的关键字参数
client.listen(...)的关键字参数直接映射到线上的SubscriptionFilter(在 src/mcp/client/client.py 中定义),支持四种订阅类型:
| 参数 | 类型 | 含义 |
|---|---|---|
tools_list_changed | bool | 订阅工具列表变化(ToolsListChanged) |
prompts_list_changed | bool | 订阅提示词列表变化(PromptsListChanged) |
resources_list_changed | bool | 订阅资源列表变化(ResourcesListChanged) |
resource_subscriptions | Sequence[str] | 订阅一组资源 URI 的内容变化(ResourceUpdated(uri=...)) |
注意resource_subscriptions接收的是URI 序列。如果误传一个裸字符串,SDK 在 listen 实现 中会直接抛出TypeError提醒你。
四个类型化事件
迭代产生四种类型化事件,它们定义在 src/mcp/shared/subscriptions.py,服务端与客户端共用:
ToolsListChanged—— 工具列表变了PromptsListChanged—— 提示词列表变了ResourcesListChanged—— 资源列表变了ResourceUpdated(uri=...)—— 某个 URI 对应的资源内容变了
事件只告诉你"什么"变了,从不告诉你"怎么"变的。这正是follow_board在收到事件后主动调用read_resource和list_tools的原因:事件只是"重新获取数据"的提示信号(cue),永远不是数据本身(payload)。这也是整个订阅机制最重要的心智模型——两端(客户端与服务端)都是"重新拉取"而非"推送内容"。
读取事件 URI 而不是自行假设
当过滤器里包含多个 URI 时,服务端可能报告的是其中某个 URI 的子资源(sub-resource)发生了变化(这是协议规范允许的)。因此不要假设"过滤器里只有一个 URI,变了的就是它",而应直接读取event.uri来判断。这与服务端MCPServer的"精确字符串匹配"行为形成对照:MCPServer只对完全匹配的 URI 发布通知,但客户端必须做好收到子资源 URI 事件的准备。
重复事件的合并
处于"待消费"状态(unconsumed)的完全相同的重复事件会被合并为一个:多个ToolsListChanged还没被消费时,只会在队列里留一个。合并的前提是"完全相同"——两个指向不同 URI 的ResourceUpdated就是两个独立事件。这一设计的意义在于:事件是"级别触发"(level trigger)而非"边沿触发",重新获取数据总会拿到当前最新状态,所以合并并不会造成信息丢失,反而能在高并发下显著减少重复的重新拉取。
订阅句柄的两个重要属性
listen()返回的订阅对象(Subscription)上有两个值得关注的属性:
sub.honored:服务器最终确认(acknowledge)的过滤器,是一个SubscriptionFilter对象,可以用属性访问的方式读取你传入的字段,例如sub.honored.prompts_list_changed。MCPServer会 honor 你请求的每一种类型,因此通常会把你的请求原样回显;而能力较少的服务器只会确认较少的部分,并且被 honor 的类型也未必真的会触发。服务器也可能拒绝整个请求(而不是确认),这会在请求层面表现为错误——服务端如何决定"谁可以看",见 docs/handlers/subscriptions.md。sub.subscription_id:该 listen 请求的 JSON-RPC id,也是这条流上每一帧(frame)都携带的订阅 id。在同一客户端上可以同时打开多条订阅,每条流凭自己的 id 被多路复用(demultiplex)区分。在 python-sdk 中,客户端使用"listen-1"、"listen-2"这样的字符串 id(由进程级计数器生成,见 src/mcp/client/subscriptions.py),而其他客户端可能使用整数 id。
不阻塞主流程的并行监听
follow_board会一直运行到服务器关闭流——而服务器可能永远不会主动关闭。因此如果单独运行,它会霸占整个程序。真实世界的客户端需要 watcher 与主流程并行:agent 继续调用工具,同时 watcher 持续刷新缓存或 UI。
正确顺序是:先打开订阅,再启动 watcher,然后继续干自己的活。SDK 为三大异步生态分别提供了等价的示例,以下三个文件在 docs_src/subscriptions/ 目录下:
=== "asyncio"
```python title="app.py" import asyncio from mcp import Client from mcp.client.subscriptions import Subscription from .tutorial003 import BOARD, read_board async def watch(client: Client, sub: Subscription) -> None: async for _event in sub: board = await read_board(client) print(board) if "[ ]" not in board: return # sprint finished: the stream closes when run_sprint leaves the block async def run_sprint(client: Client) -> None: async with client.listen(resource_subscriptions=[BOARD]) as sub: print(await read_board(client)) # snapshot: acknowledged, so nothing after this is missed watcher = asyncio.create_task(watch(client, sub)) for task in ("design", "build", "ship"): await client.call_tool("complete_task", {"board": "sprint", "task": task}) await watcher # returns once the watcher has seen the finished board async def main() -> None: async with Client("http://localhost:8000/mcp") as client: await run_sprint(client) if __name__ == "__main__": asyncio.run(main()) ```=== "trio"
```python title="app.py" import trio from mcp import Client from mcp.client.subscriptions import Subscription from .tutorial003 import BOARD, read_board async def watch(client: Client, sub: Subscription) -> None: async for _event in sub: board = await read_board(client) print(board) if "[ ]" not in board: return # sprint finished: the stream closes when run_sprint leaves the block async def run_sprint(client: Client) -> None: async with client.listen(resource_subscriptions=[BOARD]) as sub: print(await read_board(client)) # snapshot: acknowledged, so nothing after this is missed async with trio.open_nursery() as nursery: nursery.start_soon(watch, client, sub) for task in ("design", "build", "ship"): await client.call_tool("complete_task", {"board": "sprint", "task": task}) async def main() -> None: async with Client("http://localhost:8000/mcp") as client: await run_sprint(client) if __name__ == "__main__": trio.run(main) ```=== "anyio"
```python title="app.py" import anyio from mcp import Client from mcp.client.subscriptions import Subscription from .tutorial003 import BOARD, read_board async def watch(client: Client, sub: Subscription) -> None: async for _event in sub: board = await read_board(client) print(board) if "[ ]" not in board: return # sprint finished: the stream closes when run_sprint leaves the block async def run_sprint(client: Client) -> None: async with client.listen(resource_subscriptions=[BOARD]) as sub: print(await read_board(client)) # snapshot: acknowledged, so nothing after this is missed async with anyio.create_task_group() as tg: tg.start_soon(watch, client, sub) for task in ("design", "build", "ship"): await client.call_tool("complete_task", {"board": "sprint", "task": task}) async def main() -> None: async with Client("http://localhost:8000/mcp") as client: await run_sprint(client) if __name__ == "__main__": anyio.run(main) ```(上述app.py示例从第一个示例中导入BOARD和read_board,仓库中将其保存为tutorial003.py。如果你把渲染后的文件分别保存为client.py和app.py,则应改为from client import BOARD, read_board;下文watch.py示例同样以相同方式导入read_board。)
顺序就是一切
三个示例的核心逻辑完全一致,关键洞察有两点:
- 没有任何重放(replay)。流创建之前发布的事件会永远错过。而
client.listen(...)的进入会等待服务器确认,因此从确认那一刻起的每一个变化都会到达你的 watcher——在 block 内部取的快照(snapshot)不会漏掉任何一个变化。所以顺序必须是:打开订阅 → 确认完成 → 拍快照 → 启动 watcher。 - 同一 client 上的其他请求与流完全并行。无论是 watcher 任务发起的还是其他任务发起的请求,都可以在同一条打开的流旁边自由运行。由于"未消费的重复事件会合并",繁忙的主流程可能只需一次重新拉取(refetch)而不是三次;而不同的事件不会合并——命名多个 URI 的过滤器会为每个 URI 各自维护一个待处理事件队列。
如何停止监听:退出 block 就是取消订阅
停止监听的唯一方式是退出上下文管理器 block——没有unsubscribe()这样的调用。取消拥有该 block 的任务会自动完成这一点:SDK 会按照传输层(transport)期望的方式取消 listen 请求,例如在 Streamable HTTP 传输上,就是关闭该请求的流。如果一个 watcher 要存活整个应用的生命周期,它永远不会自行返回,因此在应用关闭(shutdown)时,需要显式取消它本身或其所属的 task group 的 scope。
流的结束:两种结局,一样的对策
流只有两种结束方式,而两者都是普通的控制流(control flow):
- 优雅关闭(graceful close):服务器主动关闭流,
async for循环自然结束。 - 突然中断(abrupt drop):连接意外断开,循环抛出
SubscriptionLost。
这个区别只用于诊断,并不改变接下来的行动——反正流已经没了、什么都不会重放,还在乎的 watcher 需要重新 listen 并重新拉取数据。示例代码来自 docs_src/subscriptions/tutorial005.py:
import anyio from mcp import Client from mcp.client.subscriptions import SubscriptionLost from .tutorial003 import read_board async def keep_following(client: Client) -> None: while True: try: async with client.listen(resource_subscriptions=["board://sprint"]) as sub: print(await read_board(client)) # refetch: no replay across streams async for _event in sub: print(await read_board(client)) except SubscriptionLost: pass # Either ending means the stream is gone. Back off before re-listening: # a graceful close may be the server shedding load. await anyio.sleep(1)优雅关闭不等于"别再订阅了"
服务器可能出于自己的原因优雅地关闭流——比如某个订阅者的积压(backlog)过大,被服务器主动"甩掉"(shed)。因此干净的结束不是停止监听的信号。keep_following在两种结束方式之后都会await anyio.sleep(1)退避(back off)再重新 listen,这是对服务器的一种基本礼貌。
本地的 SubscriptionLost 成因:1024 个未消费事件上限
SubscriptionLost还有一个客户端本地成因:客户端最多缓存 1024 个未消费事件(_MAX_PENDING_EVENTS = 1024,见 src/mcp/client/subscriptions.py)。消费速度落后到超过这个上限的消费者,会直接失去订阅,而不是让内存无限增长。这提示了一个重要的编码习惯:保持async for的循环体短小精悍,把耗时的工作放到循环外面去做(例如只做入队或通知,真正的重活交给其他任务)。
进入 listen() 可能抛出的异常
keep_following只捕获了SubscriptionLost,但进入listen()时还可能抛出其他异常(见 listen 的 docstring 与 tests/client/test_subscriptions.py 中的行为测试):
| 异常 | 触发条件 | 是否值得重试 |
|---|---|---|
MCPError | 连接失败,或服务器不提供该方法(未注册 listen 处理) | 视情况 |
TimeoutError | 在会话读超时时间内没有收到服务器确认 | 通常值得 |
ListenNotSupportedError | 连接协商出的协议版本早于 2026(不支持subscriptions/listen) | 永不——重试也不会好转,应改用旧版subscribe_resource()路径 |
SubscriptionLost | 流在确认之前就结束了 | 值得,配合退避 |
需要你自行决定 watcher 对其中哪些异常进行重试;其中最后一个(ListenNotSupportedError)永远不会自我修复。
源码视角:listen 的底层机制
理解了使用层面,再来看 SDK 内部是如何实现"进入即确认、迭代即消费、退出即取消"这条契约的。整个驱动在 src/mcp/client/subscriptions.py 中,由三个部件协作:
1.ListenRoute:一条流的路由与去重状态
每条 listen 流对应一个ListenRoute对象(src/mcp/client/subscriptions.py#L75-L148),由会话(session)在收到确认前就预先注册(_register_listen_route),从而保证"确认与流上帧的到达不会竞争"。它维护:
honored:服务器确认的过滤器;acked:确认到达的anyio.Event,listen()的进入等待的就是它;_pending:以待消费事件为键的字典——键就是去重的手段,相同事件入队时被字典天然吸收,这正是"重复未消费事件合并"的底层实现;_honored_uris:被 honor 的资源 URI 集合,用于判定ResourceUpdated事件是否在订阅范围内。
值得注意的是,deliver()对ResourceUpdated的准入判断是"只要 URI 订阅被 honor 就放行"——因为协议允许事件携带的 URI 是被订阅 URI 的子资源,无法提前精确匹配。事件队列长度触及_MAX_PENDING_EVENTS(1024)时,route 以"lost"结局收场,并附带一条说明积压超限的错误。
2.Subscription:暴露给用户的异步迭代器
Subscription对象(src/mcp/client/subscriptions.py#L155-L197)是对ListenRoute的薄封装。__anext__调用route.next_event():拿到事件就返回;拿到"lost"结局就抛出SubscriptionLost(并把底层错误链在 cause 上);拿到优雅结局就抛出StopAsyncIteration结束循环。实现细节上,next_event会先"快照"唤醒事件再检查状态,保证事件送达不会与检查竞争而丢失;"local"(本地退出)结局会直接短路积压队列,而优雅结束等其他结局会先排空积压——一次优雅关闭绝不会吞掉它之前已到达的事件。
3.listen:异步上下文管理器
listen()(src/mcp/client/subscriptions.py#L200-L282)是@asynccontextmanager:
- 进入:检查协议版本(低于 2026-07-28 抛
ListenNotSupportedError)→ 构造SubscriptionsListenRequest→ 用"listen-N"格式的字符串 id 发送请求 → 在会话读超时内等待确认。 - 确认:服务器回显的过滤器被记录到
sub.honored;如果服务器直接以结果帧回应(视为"打开即已关闭"的退化场景),则 honor 一个空过滤器。 - 退出:
finally中把 route 结算为"local",取消驱动任务并注销路由——这就是"退出 block 即取消订阅"的实现。驱动的请求刻意不设结果超时,因为响应要等到流结束时才会到来。
服务端契约的呼应
这套客户端契约与服务端的MCPServer实现遥相呼应:MCPServer自动承担线上的义务——确认作为第一帧、按流过滤、订阅 id 打在每一帧上。上线帧形如{"method": "notifications/subscriptions/acknowledged", "params": {"notifications": {...}, "_meta": {"io.modelcontextprotocol/subscriptionId": "listen-1"}}},随后是notifications/resources/updated。注意更新帧不携带看板内容,只携带_meta中的订阅 id——与"事件是提示而非载荷"的设计一脉相承。服务端的完整故事(发布事件、收窄过滤器、跨进程扩展的SubscriptionBus、ListenHandler)在 docs/handlers/subscriptions.md。
要点回顾
- 进入
async with client.listen(...);进入动作会等待服务器确认,因此确认之后发布的一切都不会错过。 - 用
async for event in sub迭代。事件是重新拉取的提示,永远不是数据本身。 - 先打开订阅,再把 watcher 跑成任务,工具调用继续与它并行;进入时在 block 内拍快照,保证不漏变化。
- 优雅结束停止循环;意外中断抛出
SubscriptionLost。无论哪种:重新 listen、重新拉取、先退避。 - 退出 block 就是取消订阅,没有独立的
unsubscribe调用;应用关闭时记得取消长命 watcher 或其 task group scope。 - 保持
async for循环体短小,避免消费落后触发 1024 事件积压上限。
这些事件同样维持着客户端缓存的新鲜度——这正是下一篇 Caching 要讲的内容;而如何发布这些事件、如何收窄过滤器、如何扩展出单进程边界,请阅读服务端篇 Subscriptions。
【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址: https://gitcode.com/gh_mirrors/pythonsd/python-sdk
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考