十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

MCP Python SDK 服务端订阅机制全解析:从 `subscriptions/listen` 到跨进程扩展

MCP Python SDK 服务端订阅机制全解析:从 `subscriptions/listen` 到跨进程扩展 人工智能MCP 服务MCP Clients【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址https://gitcode.com/gh_mirrors/pythonsd/python-sdk点击查看免费下载本篇文章以 Model Context Protocol 官方 Python SDK即本仓库src/mcp中的服务端订阅功能为主线深入讲解 MCP 2026-07-28 时代SEP-2575的subscriptions/listen机制客户端如何通过一次请求拿到一条常驻的变更通知流服务端如何用一行ctx.notify_*发布变更过滤器与订阅 ID 在线上如何呈现以及如何在低层Server上手动组装、如何用SubscriptionBus把订阅扩展到多副本部署。读完本文你将能够从零实现一个带实时资源/工具变更通知的 MCP 服务器并理解其背后的 SDK 实现细节。为什么需要订阅服务器的目录不是固定的服务器的目录catalog并不是一成不变的工具可能在运行时才出现例如某个功能开关开启后才注册新工具资源 URI 背后的内容也会随时变化。客户端如何得知这些变化答案就是订阅Subscriptions。客户端发送一次subscriptions/listen请求而这次请求的响应本身就是一条流它保持打开状态持续承载客户端所请求的变更通知直到流被关闭。这与 2025 时代的resources/subscribe走的是完全不同的路径——关于这一点下文「它不是什么」一节会详细对比。从工具中发布变更服务端只做一行服务端要做的其实只有一件事发布变更。以下面这份完整的服务端示例为准完整可运行代码见 docs_src/subscriptions/tutorial001.pyfrom mcp.server.mcpserver import Context, MCPServer mcp MCPServer(Sprint Board) BOARDS { sprint: {design: False, build: False, ship: False}, backlog: {tidy docs: False}, } mcp.resource(board://{name}) def board(name: str) - str: tasks BOARDS[name] return \n.join(f[{x if done else }] {task} for task, done in tasks.items()) mcp.tool() async def complete_task(board: str, task: str, ctx: Context) - str: BOARDS[board][task] True await ctx.notify_resource_updated(fboard://{board}) return f{task}: done def sprint_report() - str: done sum(done for tasks in BOARDS.values() for done in tasks.values()) return f{done} task(s) done mcp.tool() async def enable_reports(ctx: Context) - str: mcp.add_tool(sprint_report) await ctx.notify_tools_changed() return reporting is live这里有四条发布 API覆盖了四种可能变化的目录await ctx.notify_resource_updated(board://sprint)通知该 URI 内容已变化。它只送达订阅了这个 URI的打开中的流其他流一概收不到。await ctx.notify_tools_changed()通知工具列表已变化。收到此通知的客户端会再次调用tools/list这次就能看到新注册的sprint_report工具了。它的两个兄弟方法notify_prompts_changed()与notify_resources_changed()分别对应提示词prompts列表与资源resources列表的变化。注意一个关键设计没有订阅者就没有任何工作。向空闲的服务器发布变更是一个 no-op空操作——你永远不需要去检查有没有人在听你只需要声明「什么变了」。这大大简化了服务端业务代码。MCPServer替你处理了subscriptions/listen请求。线上的通信义务——首帧的确认应答acknowledgment、按流过滤、给每一帧打上订阅 ID——全部是 SDK 的责任服务端业务代码无需关心。线上长什么样确认帧、订阅 ID 与通知帧当complete_task执行后一条过滤器指定了board://sprint的流在线上会看到这样的帧序列{method: notifications/subscriptions/acknowledged, params: {notifications: {resourceSubscriptions: [board://sprint]}, _meta: {io.modelcontextprotocol/subscriptionId: listen-1}}} {method: notifications/resources/updated, params: {uri: board://sprint, _meta: {io.modelcontextprotocol/subscriptionId: listen-1}}}请注意这个更新没有携带的东西它不携带面板内容本身。它只是一个「变了」的信号读取内容要靠客户端重新resources/read。再注意每一帧都带着 listen 请求的 JSON-RPC id它放在_meta键io.modelcontextprotocol/subscriptionId之下这个 id 就是订阅 ID。这个 ID 由客户端铸造Python 的Client使用listen-1这样的字符串其他客户端也可能用整数。SDK 实现里的细节印证从源码看订阅 ID 的常量定义在 src/mcp/shared/subscriptions.pySUBSCRIPTION_ID_META_KEY io.modelcontextprotocol/subscriptionId而线上帧的构造由event_to_notification()完成——服务端方向的每个事件都会被转换成带_meta戳的ServerNotification。帧的类型来自LISTEN_STREAM_METHODS这个冻结集合它只包含四条通知方法LISTEN_STREAM_METHODS: frozenset[str] frozenset({*_LIST_CHANGED_EVENTS, notifications/resources/updated})即notifications/tools/list_changed、notifications/prompts/list_changed、notifications/resources/list_changed与notifications/resources/updated——这就是 2026-07-28 时代能骑在 listen 流上的全部变更通知词汇表。只送达被要求的内容过滤器是契约过滤器是一份契约。一条同时请求了「工具列表变化」和「一个资源 URI」的流只会收到这两种事件其他一律收不到。哪怕服务端发布了提示词变更那条流也保持沉默。具体的匹配规则如下MCPServer对资源 URI 做精确字符串匹配。指定了board://sprint的流听不到任何关于board://sprint/tasks/1的消息。规范允许服务器报告「已订阅 URI 的子资源」的变化MCPServer不会这么做但客户端是按这个预期来构建的——也就是说你的客户端代码不应该假设收不到子资源通知。这条「双方共享的准入谓词」在源码里有明确的实现即 src/mcp/shared/subscriptions.py 的event_matches()工具列表变化只匹配tools_list_changed is True资源更新事件则要判断event.uri in urishonored 的 URI 集合。服务器端的投递和客户端的接收都只认被确认acknowledged的过滤器子集两端是同一套逻辑。另外在服务端 src/mcp/server/subscriptions.py 中_honored_subset()会构造「服务器将投递的子集」用于确认帧所有被请求的种类都会被 honored——某类事件是否真的触发取决于服务器实际发布了什么这就像订阅一个不存在的资源 URI 一样被 honored但永不触发。这条流不是什么有两个容易误解的点需要澄清它不是重放日志replay log。断掉的流就是断了没有人连接期间发布的事件不会被排队。客户端需要重新 listen 并重新拉取re-listen and refetch。在 src/mcp/client/subscriptions.py 中Subscription.__anext__在流未经服务端优雅关闭而中断时会抛出SubscriptionLost提示信息正是 re-listen and refetch。它不是 2025 时代的路径。调用过resources/subscribe的客户端由ctx.session.send_resource_updated(uri)服务notify_*方法只会到达subscriptions/listen的流。两者是完全独立的机制。决定谁能监视用中间件做访问门控默认情况下每一种被请求的类型和 URI 都会被接受——任何调用方都可以监视你发布的任何 URI。而且这里有一个微妙的安全缺口没有任何代码会去查你的读取处理器read handler因为根本没人读取。一个会被你的files://{name}处理器拒之门外的调用方仍然可以打开files://payroll.csv的流得知它变了、以及什么时候变的。但它永远学不到内容也无法探测存在性——因为未知的 URI 同样会被接受只是永远不会触发而已。这个缺口虽然窄却是真实存在的所以在多租户服务器上发布「按用户区分」的 URI 之前务必加上门控。门控就是中间件middleware它在 SDK 确认请求之前看到subscriptions/listen请求当调用方请求了它无权读取的内容时直接拒绝。完整示例见 docs_src/subscriptions/tutorial006.pyfrom mcp_types import INVALID_REQUEST, SubscriptionsListenRequestParams from mcp.server.auth.middleware.auth_context import get_access_token from mcp.server.context import CallNext, HandlerResult, ServerRequestContext from mcp.server.mcpserver import MCPServer from mcp.shared.exceptions import MCPError # Who may see each file. Replace this table with a database or your RBAC system. ACCESS { files://report.pdf: {alice, bob}, files://payroll.csv: {carol}, } def can_access(user: str | None, uri: str) - bool: return user is not None and user in ACCESS.get(uri, set()) async def gate_subscriptions(ctx: ServerRequestContext, call_next: CallNext) - HandlerResult: if ctx.method subscriptions/listen: params SubscriptionsListenRequestParams.model_validate(ctx.params or {}, by_nameFalse) token get_access_token() user token.subject if token else None if not all(can_access(user, uri) for uri in params.notifications.resource_subscriptions or ()): raise MCPError(INVALID_REQUEST, not permitted to watch the requested resources) return await call_next(ctx) mcp MCPServer(Reports, middleware[gate_subscriptions]) mcp.resource(files://{name}) def file(name: str) - str: uri ffiles://{name} token get_access_token() if not can_access(token.subject if token else None, uri): raise MCPError(INVALID_REQUEST, fUnknown resource: {uri}) return fcontents of {name}这套门控有四个要点ctx.params是原始请求。中间件需要自己把它校验为SubscriptionsListenRequestParams通过model_validate(..., by_nameFalse)从而读出客户端请求的过滤器。拒绝方式是在call_next(ctx)之前抛出MCPError。客户端会收到这个错误拿不到流但连接本身继续存活。错误消息要保持统一、不要点名具体 URI否则一次拒绝就会暴露哪些 URI 是受保护的。一个can_access(user, uri)同时回答两个问题资源处理器在resources/read时询问它中间件在subscriptions/listen时询问它。把示例里的静态表换成数据库或你的 RBAC 系统两端依然步调一致不会出现「能读却订不了」或「订得了却读不了」的错位。判定在流的整个生命周期内有效。没有「每个事件重新检查」的机制。所以如果调用方的访问权限可能在流中途失效例如过期的 token你必须在失效时主动结束该调用方的连接。关于中间件契约的完整说明——包括它还包装什么、为什么被标记为 provisional暂定——见 docs/advanced/middleware.md。客户端一侧async with client.listen(...)流的那一端是一个追着面板跑的客户端。完整代码见 docs_src/subscriptions/tutorial003.pyfrom 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_changedTrue, resource_subscriptions[BOARD]) as sub: async for event in sub: match event: case ResourceUpdated(uriuri): 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(...)时请求就已发出并且会等待服务器的确认应答——所以当async with代码块开始执行时流已经是活的了。每一个类型化事件都只是「重新拉取」的提示cue绝不是载荷。收到ResourceUpdated就重新read_resource收到ToolsListChanged就重新list_tools。过滤器没请求的种类永远不会到达——所以case _:分支在正常流程中不会命中。从客户端源码看src/mcp/client/subscriptions.pylisten()是一个异步上下文管理器进入时发送SubscriptionsListenRequest并等待确认退出时结束订阅。它会在以下情形抛出明确异常协商的协议版本早于 2026-07-28 时抛ListenNotSupportedError服务器拒绝请求或连接在确认前失败时抛MCPError流在确认前就结束时抛SubscriptionLost确认前读取超时抛TimeoutError。另外resource_subscriptions参数要求传入 URI 序列——传入裸字符串会直接抛TypeError。listen()还支持on_event回调客户端用它来完成缓存驱逐确保消费者重新拉取前缓存已失效。关于客户端的其余内容——如何与主流程并行监视、流的结束方式、如何重新 listen——有专门的页面见 docs/client/subscriptions.md。跨进程扩展实现你自己的SubscriptionBus发布从你的处理器到打开的流中间经过一个SubscriptionBus。默认实现是内存版一个进程、进程内所有流。在跑负载均衡器后面的副本之前这个默认实现就是正确答案。但一旦跑起多副本情况就变了客户端的流被钉在某一个副本上而另一个副本上的发布必须能到达它。这个接缝seam需要你自己实现在你的 pub/sub 后端之上实现两个方法即可。从源码看SubscriptionBus是一个Protocol见 src/mcp/server/subscriptions.py只要求两个方法class SubscriptionBus(Protocol): async def publish(self, event: ServerEvent) - None: ... def subscribe(self, listener: Callable[[ServerEvent], None]) - Callable[[], None]: ...publish是异步的因为后端实现要做网络 I/Osubscribe是同步的本地注册返回一个幂等的 unsubscribe 可调用对象。监听器是同步函数、不得抛异常、运行在服务器的事件循环上。下面是文档给出的基于 Redis 的示例实现from collections.abc import Callable from redis.asyncio import Redis from mcp.server.mcpserver import MCPServer from mcp.server.subscriptions import ServerEvent # SubscriptionBus is a Protocol: no base class class RedisSubscriptionBus: def __init__(self, redis: Redis) - None: self._redis redis self._listeners: dict[object, Callable[[ServerEvent], None]] {} async def publish(self, event: ServerEvent) - None: await self._redis.publish(mcp-events, encode(event)) # to every replica def subscribe(self, listener: Callable[[ServerEvent], None]) - Callable[[], None]: token object() self._listeners[token] listener def unsubscribe() - None: self._listeners.pop(token, None) return unsubscribe mcp MCPServer(Sprint Board, subscriptionsRedisSubscriptionBus(redis))encode由你自己实现每个副本上还需要一个读任务reader task负责解码到达的消息并调用所有已注册的监听器。注意这里的令牌机制InMemorySubscriptionBus与示例都用object()作为字典键这样同一个 callable 可以注册多次因为绑定方法之间比较相等。两个重要的架构边界总线搬运的是类型化的ServerEvent值而不是 JSON-RPC。事件词汇表是四个小数据类定义在 src/mcp/shared/subscriptions.pyToolsListChanged、PromptsListChanged、ResourcesListChanged三者都是无字段的冻结数据类和携带uri的ResourceUpdated。打订阅 ID、过滤、流生命周期都留在 SDK 里所以一个总线实现不可能破坏协议——它唯一能做的就是跨进程搬运事件。这正是把总线设计成 Protocol 而不是基类的意义所在源码注释明确写着 SubscriptionBus is a Protocol: no base class。另外内存版InMemorySubscriptionBus.publish会逐个调用监听器并做「扇出边界隔离」某个监听器抛异常会被记录并跳过不会饿死其他监听器也不会让发布方处理器失败最后还有一个anyio.lowlevel.checkpoint()让单任务里的一连串发布有机会让 listen 流在两个事件之间排空而不是溢出不读的缓冲区。在请求之外发布自己持有总线引用想从请求之外发布例如 lifespan 任务、webhook就要自己构建总线以持有引用。因为当你什么都不传时MCPServer会在内部建一个总线但不会把它暴露出来from mcp.server.subscriptions import InMemorySubscriptionBus, ToolsListChanged bus InMemorySubscriptionBus() mcp MCPServer(Sprint Board, subscriptionsbus) async def tools_reloaded() - None: await bus.publish(ToolsListChanged()) # from a lifespan task, a webhook, anywhere低层组装Server上的三行手工接线在低层的Server上没有任何预接线的部件同样的零件三行就能组装起来。完整示例见 docs_src/subscriptions/tutorial002.pyfrom typing import Any import mcp.types as types from mcp.server.context import ServerRequestContext from mcp.server.lowlevel import Server from mcp.server.subscriptions import InMemorySubscriptionBus, ListenHandler, ResourceUpdated bus InMemorySubscriptionBus() listen_handler ListenHandler(bus) BOARD {design: False, build: False} COMPLETE_TASK_SCHEMA: dict[str, Any] { type: object, properties: {task: {type: string}}, required: [task], } async def read_resource( ctx: ServerRequestContext[Any], params: types.ReadResourceRequestParams ) - types.ReadResourceResult: board \n.join(f[{x if done else }] {task} for task, done in BOARD.items()) return types.ReadResourceResult(contents[types.TextResourceContents(uriparams.uri, textboard)]) async def list_tools( ctx: ServerRequestContext[Any], params: types.PaginatedRequestParams | None ) - types.ListToolsResult: return types.ListToolsResult( tools[types.Tool(namecomplete_task, descriptionMark a task done., input_schemaCOMPLETE_TASK_SCHEMA)] ) async def call_tool(ctx: ServerRequestContext[Any], params: types.CallToolRequestParams) - types.CallToolResult: args params.arguments or {} BOARD[args[task]] True await bus.publish(ResourceUpdated(uriboard://sprint)) return types.CallToolResult(content[types.TextContent(typetext, textdone)]) server Server( sprint-board, on_read_resourceread_resource, on_list_toolslist_tools, on_call_toolcall_tool, on_subscriptions_listenlisten_handler, )三个要点总线归你所有所以你可以直接向它发布await bus.publish(ResourceUpdated(uri...))。把它放在你的处理器能到达的地方——这里放在模块作用域更大的应用则放在 lifespan 里。ListenHandler(bus)正是MCPServer注册的同一个处理器而on_subscriptions_listen是一个普通处理器槽位。如果你想要不同的语义可以在该槽位放入你自己的 callable——但那样规范义务就转移给你了先确认、给每一帧打订阅 ID、绝不投递过滤器之外的内容。ListenHandler.close()会优雅地结束所有打开的流。每条流都会把 listen 请求的 result 作为最后一帧收到——这是规范中「服务器有意结束订阅」的表示方式。注意它在流刷完之前就会返回所以在拆除传输层之前要给流留一点时间如果不调用它流会在客户端断开时自然结束。从 src/mcp/server/subscriptions.py 看ListenHandler还有两个可调参数用于保护服务器资源max_subscriptions默认 1024限制并发流的数量超过后新的 listen 请求会在确认前以INTERNAL_ERROR被拒绝max_buffered_events默认 1024限制每条流的待投递事件积压积压达到上限的流会被结束客户端重新 listen 并 refetch——由于没有重放结束流不会丢失积压之外本就会丢失的东西。缓冲的设计还保证发布方不会被慢消费者阻塞。此外ListenHandler内部「先订阅总线、再发送确认帧」的顺序也值得注意这样在确认帧写入被挂起期间发布的事件会被缓冲而不是丢失而确认帧依然保证是首帧因为只有这一个任务写流且它只在确认发送返回后才开始排空缓冲区。总结客户端用一次subscriptions/listen请求完成订阅响应本身就是流。服务端处理它是内置能力。你用ctx.notify_*发布变更打订阅 ID、过滤、生命周期处理全部由 SDK 负责。事件是提示不是载荷。两端都是「重新拉取」模式。客户端一侧是async with client.listen(...)完整故事见 docs/client/subscriptions.md。在低层Server上你自己组装同样的零件一个总线、ListenHandler(bus)、on_subscriptions_listen槽位。扩展scale out就是实现SubscriptionBus两个方法并以MCPServer(subscriptions...)传入。默认的InMemorySubscriptionBus是单进程的正确答案多副本时用 Redis/NATS 之类的 pub/sub 后端实现同一接口即可事件词汇表是四个类型化数据类协议安全由 SDK 保证。关于承载这一切的服务器如何在单副本到二十副本的规模下部署运行见 docs/run/deploy.md。对应的测试用例可以在 tests/server/test_subscriptions.py、tests/client/test_subscriptions.py 以及 tests/docs_src/test_subscriptions.py 中找到是理解本文所述行为的可运行佐证。赞分享人工智能MCP 服务MCP Clients【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址https://gitcode.com/gh_mirrors/pythonsd/python-sdk点击查看免费下载相关推荐python-sdk 服务器端订阅机制全解析从 subscriptions/listen 流到多副本横向扩展python sdk 服务器端订阅机制全解析从 subscriptions/listen 流到多副本横向扩展 Model Context ProtocolM人工智能MCP 服务MCP ClientsMCP Python SDK 服务端订阅机制实战subscriptions/listen 事件流、过滤器与多进程扩展MCP Python SDK 服务端订阅机制实战subscriptions/listen 事件流、过滤器与多进程扩展 服务器目录并非一成不变工具会在运行时出人工智能MCP 服务MCP ClientsMCP Python SDK 订阅机制全解析从 subscriptions/listen 流式通知到多副本扩展MCP Python SDK 订阅机制全解析从 subscriptions/listen 流式通知到多副本扩展 服务器的目录并非一成不变工具会在运行时出现人工智能MCP 服务MCP Clients上一篇重新定义浏览器自动化基于MCP协议的智能代理架构革命下一篇flags开源功能标识库助力项目快速开发创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表