news 2026/9/20 10:29:21

python-sdk 客户端订阅完全指南:用 client.listen() 实时监听 MCP 资源与目录变化

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
python-sdk 客户端订阅完全指南:用 client.listen() 实时监听 MCP 资源与目录变化

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_changedbool订阅工具列表变化(ToolsListChanged
prompts_list_changedbool订阅提示词列表变化(PromptsListChanged
resources_list_changedbool订阅资源列表变化(ResourcesListChanged
resource_subscriptionsSequence[str]订阅一组资源 URI 的内容变化(ResourceUpdated(uri=...)

注意resource_subscriptions接收的是URI 序列。如果误传一个裸字符串,SDK 在 listen 实现 中会直接抛出TypeError提醒你。

四个类型化事件

迭代产生四种类型化事件,它们定义在 src/mcp/shared/subscriptions.py,服务端与客户端共用:

  • ToolsListChanged—— 工具列表变了
  • PromptsListChanged—— 提示词列表变了
  • ResourcesListChanged—— 资源列表变了
  • ResourceUpdated(uri=...)—— 某个 URI 对应的资源内容变了

事件只告诉你"什么"变了,从不告诉你"怎么"变的。这正是follow_board在收到事件后主动调用read_resourcelist_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_changedMCPServer会 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示例从第一个示例中导入BOARDread_board,仓库中将其保存为tutorial003.py。如果你把渲染后的文件分别保存为client.pyapp.py,则应改为from client import BOARD, read_board;下文watch.py示例同样以相同方式导入read_board。)

顺序就是一切

三个示例的核心逻辑完全一致,关键洞察有两点:

  1. 没有任何重放(replay)。流创建之前发布的事件会永远错过。而client.listen(...)的进入会等待服务器确认,因此从确认那一刻起的每一个变化都会到达你的 watcher——在 block 内部取的快照(snapshot)不会漏掉任何一个变化。所以顺序必须是:打开订阅 → 确认完成 → 拍快照 → 启动 watcher。
  2. 同一 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.Eventlisten()的进入等待的就是它;
  • _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——与"事件是提示而非载荷"的设计一脉相承。服务端的完整故事(发布事件、收窄过滤器、跨进程扩展的SubscriptionBusListenHandler)在 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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/20 10:27:14

通道剪枝三范式:Slimming、L1-norm与AutoSlim原理与实战

简介:本资源是一套面向AI算法工程师与深度学习研究者的模型轻量化实践代码包,聚焦L1-norm剪枝、Slimming通道剪枝及AutoSlim自动化结构压缩三大主流技术,解决大模型部署中参数量高、推理延迟大、硬件资源受限等实际问题。压缩包共32个文件&am…

作者头像 李华
网站建设 2026/9/20 10:26:43

转录组PCA分析三大误区:标准化、离群样本与批次效应处理

1. 先搞清楚:PCA在转录组分析里到底扮演什么角色1.1 PCA算的到底是什么做转录组分析的同学,十有八九都画过那张经典的PCA图——样本在二维平面上一颗一颗散开,组间分开了就长舒一口气,分不开就开始焦虑,甚至怀疑自己整…

作者头像 李华
网站建设 2026/9/20 10:26:37

上睑下垂分类与分割:医学图像数据集构建与双任务训练实践

简介:这套深度学习数据集聚焦眼睛及虹膜区域的上睑下垂疾病分类与分割,适用于医学图像处理、计算机视觉方向的研究者与开发者,可帮助快速搭建医疗影像分析实验。所有图像来自真实采集并自行标注,类别涵盖轻度、中度、重度、正常四…

作者头像 李华
网站建设 2026/9/20 10:24:09

Ubuntu本地部署CodeX CLI:从安装到模型对接的完整指南

直接以正文开始:如果你在Ubuntu上折腾过几款AI编程助手,大概率会有同感:网页版对话式AI和真正嵌进开发流程的编程助手,完全不是一回事。前者是“问一句答一句”,后者是“你在编辑器里写代码,它在一旁读上下…

作者头像 李华
网站建设 2026/9/20 10:21:46

SpringBoot图书借阅系统高校实战指南

简介:这是一套基于Spring Boot开发的图书借阅管理系统完整毕业设计项目,面向计算机专业本科生及Java初学者,解决高校或小型图书馆场景下的图书登记、用户管理、借阅归还、库存统计等核心业务需求。资源包共133个文件,涵盖22个Java…

作者头像 李华