YAOTU INSIGHTS

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点击查看免费下载本文基于 python-sdkModel Context Protocol 官方 Python SDK深入讲解服务器端订阅Subscriptions机制客户端如何通过一次subscriptions/listen请求获得一条常驻的事件流服务器如何在工具或资源变更时发布通知、如何用 middleware 控制谁有资格监听、如何把发布总线SubscriptionBus扩展到多进程/多副本以及低层Server上如何手工组装同一套部件。读完你将掌握在 handler 内发布变更、在客户端消费事件、并为多副本部署实现自定义事件总线的完整方案。为什么需要订阅目录不是一成不变的服务器对外公布的目录并不是静态的工具可能在运行期动态出现资源 URI 背后的内容也会变化。订阅Subscriptions就是客户端获知这些变化的方式客户端发送一个subscriptions/listen请求而该请求的响应本身就是一条流——它保持打开不断承载客户端请求过的各类变更通知。在 2026-07-28 协议SEP-2575线路上不再有常驻的 GET 流客户端通过发送subscriptions/listen主动订阅服务器事件响应即流。这条语义可以从 src/mcp/server/subscriptions.py 的模块 docstring 中得到印证。发布变更工具侧只需一行代码在服务器端你的工作只有一行把发生了什么变化发布出去。完整的示例见 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四个发布方法的行为各不相同它们在 src/mcp/server/mcpserver/context.py 中都有对应实现await ctx.notify_resource_updated(board://sprint)只会送达每一个订阅了该 URI 的打开流其他人收不到。await ctx.notify_tools_changed()送达所有请求了工具列表变化的流。收到它的客户端会再次调用tools/list此时就能看到新出现的sprint_report工具。两个同级方法notify_prompts_changed()与notify_resources_changed()分别对应提示词列表和资源列表的变化。没有订阅者就没有工作向空闲服务器发布是一个 no-op所以你永远不需要检查是否有人在听只需陈述什么变了即可。从实现看这四个notify_*方法本质上都是把对应的类型化事件投递到SubscriptionBus见 context.py 中的_bus.publish(...)真正负责把事件送到每条流上的是 SDK 内部的订阅总线。MCPServer替你承担了subscriptions/listen的服务器职责。协议层面的义务全部由 SDK 完成确认帧acknowledgment必须是流的第一帧按流做事件过滤每一帧上都携带订阅标识符subscription id。线上的样子确认帧与事件帧当过滤器指定了board://sprint的流在complete_task执行后线上的帧如下{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}}}注意更新帧并不携带看板内容本身。每一帧都在_meta中携带 listen 请求的 JSON-RPC id这个 id 就是订阅标识符。它由客户端铸造Python 的Client使用listen-1这类字符串其他客户端可能用整数。SDK 中SUBSCRIPTION_ID_META_KEY io.modelcontextprotocol/subscriptionId就是这段元数据的键名定义在 src/mcp/shared/subscriptions.py。过滤器就是契约只送达被请求的内容过滤器是一份合同。一条请求了工具列表变化 一个资源 URI的流只会收到这两种事件绝不会多。你发布一个提示词变化这条流保持沉默。一个容易踩坑的细节MCPServer对资源 URI 按精确字符串匹配。因此订阅了board://sprint的流听不到关于board://sprint/tasks/1的变化。规范允许服务器报告订阅 URI 的子资源的变化MCPServer从不这样做但客户端是按可能收到来构建的——所以客户端读取事件时应读取event.uri而非假设具体是哪个资源动了详见客户端文档 docs/client/subscriptions.md。底层过滤逻辑集中在 src/mcp/shared/subscriptions.py 的event_matchesToolsListChanged只在honored.tools_list_changed is True时通过ResourceUpdated只在该 URI 落在已确认的订阅 URI 集合内时通过——服务器投递与客户端接收共用同一个准入谓词。这条流不是什么它不是重放日志replay log。断掉的流就消失了断连期间发布的事件不会被排队。客户端重新 listen 并重新拉取数据。它不是 2025 时代的老路径。调用过resources/subscribe的旧客户端由ctx.session.send_resource_updated(uri)服务。notify_*方法只送达subscriptions/listen流context.py 的 docstring 明确说明了这一区别。谁能看用 middleware 把关访问权限默认情况下每一个被请求的种类和 URI 都会被满足任何调用者都可以观看你发布的任何 URI。此时没有人咨询你的读处理器——因为没有人读取。一个会被你的files://{name}处理器拒之门外的调用者仍然可以打开一条订阅files://payroll.csv的流得知文件变了以及何时变的。它永远不会得到内容也无法探测哪些 URI 存在因为未知 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是原始请求因此 middleware 需要自己把它校验成SubscriptionsListenRequestParamsmodel_validate(..., by_nameFalse)再读取客户端请求的过滤器。拒绝方式是在call_next(ctx)之前抛出MCPError客户端收到这个错误、得不到任何流而连接继续存活。务必保持错误信息统一、不提及任何 URI这样一次拒绝永远不会向调用者证实哪些 URI 是受保护的。用同一个can_access(user, uri)回答两个问题资源处理器在resources/read时调用它middleware 在subscriptions/listen时调用它。把示例里的表换成数据库或你的 RBAC 系统两个路径始终保持一致。决定在流的整个生命周期内生效没有逐事件的复查。如果调用者的访问权可能在流中途失效例如即将过期的 token请在失效时主动断开该客户端的连接。middleware 的完整契约包括它还包裹什么、为什么标记为 provisional见 docs/advanced/middleware.md。客户端那一端事件是重新拉取的信号下面是这条流另一端跟随看板的客户端完整代码见 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(...)时会发出请求并等待服务器的确认因此当代码块开始时流已经处于活跃状态每一个类型化事件都是重新拉取数据的信号而不是携带数据的负载。这就是整个契约的完整展示。客户端侧的更多内容如何让观察者与主流程并行、流的结束与重新监听等在独立的 docs/client/subscriptions.md 页面中迭代产出四种类型化事件ToolsListChanged、PromptsListChanged、ResourcesListChanged、ResourceUpdated(uri...)重复的未消费事件会合并句柄提供sub.honored服务器确认的过滤器与sub.subscription_idlisten 请求的 id用于多路复用多条并发订阅等属性。跨进程扩展实现你自己的SubscriptionBus发布从你的 handler 到打开的流走的是一条SubscriptionBus。默认的总线在内存中一个进程、该进程内的所有流。在负载均衡器后面跑多副本之前这是完全正确的答案——因为一旦多副本客户端的流被固定在某一个副本上而另一个副本上的发布必须能够到达它。这个接缝由你来实现在你的 pub/sub 后端之上实现两个方法。示例以 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是你的函数每个副本上负责解码到达消息并调用每个已注册 listener 的 reader 任务也是你的。Listener 是同步的、不得抛出异常、运行在服务器的事件循环上。总线搬运的是类型化的ServerEvent值——四个小型 dataclass永远不会是 JSON-RPC。打标stamping、过滤、流的生命周期都留在 SDK 里因此你的总线实现不可能破坏协议它只能把事件在进程间搬来搬去。SubscriptionBus在源码中就是一个两方法的Protocolpublish为异步、subscribe为同步本地注册并返回幂等的 unsubscribe 可调用对象定义见 src/mcp/server/subscriptions.py其内建实现InMemorySubscriptionBus在 同一文件 中发布时逐个调用 listener单个 listener 抛异常会被记录并跳过隔离扇出边界结束时还会做一次 checkpoint 让事件流有排水机会。在请求之外发布持有总线引用要想在请求之外发布例如 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同样的总线模式也出现在部署文档中无论事件由哪台服务器发布、流挂在哪个服务器对象上扇出本身不关心这些——同一进程内的两个MCPServer共享一条InMemorySubscriptionBus时行为已然如此在一个上开流、在另一个上发布流能听到。跨真实进程时 SDK 不附带任何可用的总线SubscriptionBus就是留给你在自己的 pub/sub 后端Redis、NATS 或任何你已经在跑的上实现的接缝详见 docs/run/deploy.md 的 Change notifications across replicas 一节。低层组合在没有预接线的地方自己组装在低层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...))。把它放在 handler 够得着的地方——示例放在模块级更大的应用放在 lifespan 里。ListenHandler(bus)就是MCPServer注册的同一个 handleron_subscriptions_listen只是一个普通的 handler 槽位。在这个槽位放入你自己的可调用对象以获得不同的语义届时规范义务就转移到你身上先确认、给每一帧打订阅 id、绝不投递过滤器之外的内容。ListenHandler.close()优雅地结束每一条打开的流每条流都会以 listen 请求的结果作为最后一帧收到它——这正是规范中服务器有意结束订阅的表达方式。注意该方法在流冲刷完成之前就返回所以在拆除传输之前要给它们留一点时间。不调用它流会在客户端断开时结束。ListenHandler在 src/mcp/server/subscriptions.py 中的实现还透露出几个实用细节单次调用即一条订阅流构造参数max_subscriptions1024限制并发流数超限以INTERNAL_ERROR拒绝发生在确认帧之前max_buffered_events1024限制每条流的积压事件数——当流的积压到达上限时该流会被结束客户端重新 listen 并重新拉取没有重放所以不丢任何积压之外的东西事件通过有界内存流缓冲发布者永远不会被慢消费者阻塞。订阅发生在发送确认帧之前因此确认写挂起期间发布的事件会被缓冲而不是丢失——确认帧仍是第一帧因为只有 handler 任务本身写这条流。总结客户端以一次subscriptions/listen请求加入响应即流服务它serving是内置能力。你用ctx.notify_*发布SDK 完成打标、过滤和生命周期管理。事件是信号而非负载两端都会重新拉取数据。客户端一侧就是async with client.listen(...)详见 docs/client/subscriptions.md。在低层Server上你自行组装同样的部件一条总线、ListenHandler(bus)、on_subscriptions_listen槽位。水平扩展 实现SubscriptionBus两个方法并作为MCPServer(subscriptions...)传入。在单副本或二十副本后面运行服务器的方法见 docs/run/deploy.md。赞分享人工智能MCP 服务MCP Clients【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址https://gitcode.com/gh_mirrors/pythonsd/python-sdk点击查看免费下载相关推荐MCP Python SDK 服务端订阅机制全解析从 subscriptions/listen 到跨进程扩展MCP Python SDK 服务端订阅机制全解析从 subscriptions/listen 到跨进程扩展 本篇文章以 Model Context Prot人工智能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创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考