MCP Python SDK 服务端订阅机制实战subscriptions/listen 事件流、过滤器与多进程扩展【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址: https://gitcode.com/gh_mirrors/pythonsd/python-sdk服务器目录并非一成不变工具会在运行时出现资源 URI 背后的内容也会变化。在 Model Context Protocol 的 2026-07-28 协议时代SEP-2575订阅Subscriptions是客户端获知这些变化的方式客户端发送一次subscriptions/listen请求而该请求的响应本身就是流——它保持打开持续承载客户端所请求的变更通知直到流被有意关闭。本文以官方 Python SDKmcp中 docs/handlers/subscriptions.md 为主线结合 src/mcp/server/subscriptions.py 等源码完整讲解服务端如何发布变更、如何按过滤器精确投递、如何用中间件做订阅授权以及如何通过自定义SubscriptionBus把事件投递扩展到多进程/多副本部署。读完后你将能在MCPServer与低层Server两种写法下分别实现订阅服务并掌握客户端client.listen(...)端的配合方式。从工具里发布变更一行代码服务端的职责被压缩成一句话发布变更。下面的示例来自 docs_src/subscriptions/tutorial001.py一个以board://{name}资源为中心的Sprint Board服务器from 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四个notify_*发布方法在MCPServer的请求上下文中src/mcp/server/mcpserver/context.py每个方法都只是把对应的事件发布到服务器的SubscriptionBus上await ctx.notify_resource_updated(uri)—— 通知uri对应的资源内容已变化仅送达订阅了该 URI 的流其他流不受影响await ctx.notify_tools_changed()—— 工具列表已变化收到它的客户端会重新调用tools/list从而看到新出现的sprint_reportawait ctx.notify_prompts_changed()—— 提示prompt列表已变化await ctx.notify_resources_changed()—— 资源列表已变化。没有订阅者就零开销这四个方法的底层实现都只是一次await self._bus.publish(...)。因此**没有订阅者就什么都不做**向空闲服务器发布事件是 no-op你永远不需要先检查有没有人在监听——只需要陈述什么变了。实现细节可参考 src/mcp/server/subscriptions.py 中InMemorySubscriptionBus.publish的扇出逻辑。线上的样子确认帧与事件帧MCPServer会替你实现subscriptions/listen的服务。线上的协议义务——确认帧必须是流的第一帧、按流过滤、每一帧都携带订阅 id——都由 SDK 承担。文档给出的线格式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}}}注意更新帧不携带板子的内容每一帧都在_meta下携带 listen 请求的 JSON-RPC id而这个 id 就是订阅 id。id 由客户端铸造PythonClient使用listen-1这样的字符串其他客户端可能使用整数。元数据键名定义在 src/mcp/shared/subscriptions.pySUBSCRIPTION_ID_META_KEY io.modelcontextprotocol/subscriptionId。过滤器是契约只投递被请求的内容过滤器是双方之间的契约。一个只请求了工具列表变化和某一个资源 URI 的流只会收到这两种事件其他一概不收发布一条 prompt 变更这个流保持沉默。URI 按精确字符串匹配MCPServer用event.uri in uris这种逐串比较见 src/mcp/shared/subscriptions.py 的event_matches所以订阅了board://sprint的流听不到board://sprint/tasks/1的变化。规范允许服务器报告已订阅 URI 的子资源变化MCPServer从不这么做但客户端被设计成能应对这种情况——因此客户端应读取event.uri而非假设具体是哪个资源变了。流不是什么明确这两点能避免大量误解它不是回放日志。流一旦断开就没了没有连接期间发布的事件不会被排队。客户端要重新 listen 并重新拉取。它不是 2025 年的那条老路。调用过resources/subscribe的客户端由ctx.session.send_resource_updated(uri)服务。notify_*方法只到达subscriptions/listen流。源码层面src/mcp/server/connection.py 会直接丢弃通过共享通道裸发变更通知的尝试——在这个时代变更通知只能走 listen 流。用中间件决定谁能观看默认行为是来者不拒任何请求的类别和 URI 都会被兑现任何调用者都可以观看你发布的任何 URI。没有任何东西会去咨询你的读取处理器——因为没有人真的在读取。一个本会被files://{name}处理器拒之门外的调用者仍然可以打开files://payroll.csv的流得知它变了、以及何时变的。它永远学不到内容也无法探测什么资源存在——因为未知 URI 也会被兑现只是永远不会触发事件。这种泄露很窄但真实存在因此在多租户服务器上发布按用户区分的 URI 之前务必先加闸门。闸门就是一个中间件。它在 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读取客户端请求的过滤器拒绝方式是在call_next(ctx)之前抛MCPError客户端拿到这个错误、没有流连接继续。错误信息要保持统一、不点名任何 URI这样一次拒绝永远不会确认哪些 URI 受保护一个can_access(user, uri)回答两个问题资源处理器在resources/read时问它中间件在subscriptions/listen时问它。把表换成数据库或你的 RBAC 系统两边依然同步决定在流的整个生命周期内有效没有逐事件复查。如果调用者的访问权限可能在流中途失效比如令牌过期那么失效时就应该结束该调用者的连接。中间件的完整契约包括它还包裹了什么、为什么标记为 provisional见 docs/advanced/middleware.md。客户端那一端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(...)会发送请求并等待你的确认因此块开始时流已经是活的每个键入事件都是重新拉取的提示绝不是负载。迭代产出四种类型化事件ToolsListChanged、PromptsListChanged、ResourcesListChanged、ResourceUpdated(uri...)。事件只说什么变了从不说怎么变的。客户端端的完整故事与主流程并行观看、流的结束、重新 listen在 docs/client/subscriptions.md。扩展到多进程实现 SubscriptionBus发布从处理器走向打开的流途经SubscriptionBus。默认实现是进程内的一个进程、里面的所有流。在负载均衡器后面运行副本之前这是正确的答案——因为那时一个客户端的流被钉在某一个副本上而另一个副本上的发布必须能到达它。这个接缝留给你实现在你的 pub/sub 后端之上实现两个方法。SubscriptionBus是一个Protocol见 src/mcp/server/subscriptions.py没有基类实现它即可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 任务也一样。监听器是同步的、不得抛异常、运行在服务器的事件循环上。从源码看publish是异步的以便后端实现做网络 I/Osubscribe是同步的本地注册src/mcp/server/subscriptions.py。总线承载类型化的ServerEvent值——四个小的 dataclassToolsListChanged、PromptsListChanged、ResourcesListChanged、ResourceUpdated定义于 src/mcp/shared/subscriptions.py绝无 JSON-RPC。盖章、过滤、流生命周期都留在 SDK 内所以总线实现不可能破坏协议——它只能在进程之间搬运事件。在请求之外发布自己持有总线要在请求之外发布请自己构造总线以便持有引用。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这让你可以从 lifespan 任务、webhook 等任意位置发布。低层组合on_subscriptions_listen 槽位下到低层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是一个普通处理器槽位在 src/mcp/server/lowlevel/server.py 等处声明。如果你想要不同的语义把自己的 callable 放进这个槽位那么协议义务就落到你头上先确认、给每帧盖订阅 id、过滤器之外的一律不投递ListenHandler.close()优雅地结束每个打开的流。每个流都会收到 listen 请求的结果作为最后一帧——这是规范表示服务器有意结束订阅的方式。它返回时这些流可能还没冲刷完所以在拆除传输之前给它们一点时间。没有它流会在客户端断开时结束。ListenHandler 的内置防护从源码src/mcp/server/subscriptions.py可以看到ListenHandler自带两个防护参数max_subscriptions: int 1024——并发流的上限超出后新的 listen 请求在确认之前以INTERNAL_ERROR拒绝max_buffered_events: int 1024——每条流事件积压的上限积压触顶的流会被结束客户端重新 listen 并重新拉取——既然没有回放结束流并不会比积压本身丢失更多。客户端一侧同样有 1024 个未消费事件的上限消费跟不上的订阅者会失去订阅。总结客户端用一次subscriptions/listen请求选择订阅响应就是流。提供服务是内置的。你用ctx.notify_*发布盖章、过滤、生命周期工作都由 SDK 完成。事件是提示不是负载。两端都要重新拉取。客户端端是async with client.listen(...)完整故事见 docs/client/subscriptions.md。在低层Server上你自己组装同样的部件一个总线、ListenHandler(bus)、on_subscriptions_listen槽位。扩展规模意味着实现SubscriptionBus——两个方法——并把它作为MCPServer(subscriptions...)传入。运行所有这些的服务器一个副本或二十个部署方式见 docs/run/deploy.md。【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址: https://gitcode.com/gh_mirrors/pythonsd/python-sdk创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考