:從 `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請求獲得一條常駐的事件流服務器如何在工具或資源變更時發(fā)布通知、如何用 middleware 控制誰有資格監(jiān)聽、如何把發(fā)布總線SubscriptionBus擴展到多進程/多副本以及低層Server上如何手工組裝同一套部件。讀完你將掌握在 handler 內(nèi)發(fā)布變更、在客戶端消費事件、并為多副本部署實現(xiàn)自定義事件總線的完整方案。為什么需要訂閱目錄不是一成不變的服務器對外公布的目錄并不是靜態(tài)的工具可能在運行期動態(tài)出現(xiàn)資源 URI 背后的內(nèi)容也會變化。訂閱Subscriptions就是客戶端獲知這些變化的方式客戶端發(fā)送一個subscriptions/listen請求而該請求的響應本身就是一條流——它保持打開不斷承載客戶端請求過的各類變更通知。在 2026-07-28 協(xié)議SEP-2575線路上不再有常駐的 GET 流客戶端通過發(fā)送subscriptions/listen主動訂閱服務器事件響應即流。這條語義可以從 src/mcp/server/subscriptions.py 的模塊 docstring 中得到印證。發(fā)布變更工具側只需一行代碼在服務器端你的工作只有一行把發(fā)生了什么變化發(fā)布出去。完整的示例見 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四個發(fā)布方法的行為各不相同它們在 src/mcp/server/mcpserver/context.py 中都有對應實現(xiàn)await ctx.notify_resource_updated(board://sprint)只會送達每一個訂閱了該 URI 的打開流其他人收不到。await ctx.notify_tools_changed()送達所有請求了工具列表變化的流。收到它的客戶端會再次調(diào)用tools/list此時就能看到新出現(xiàn)的sprint_report工具。兩個同級方法notify_prompts_changed()與notify_resources_changed()分別對應提示詞列表和資源列表的變化。沒有訂閱者就沒有工作向空閑服務器發(fā)布是一個 no-op所以你永遠不需要檢查是否有人在聽只需陳述什么變了即可。從實現(xiàn)看這四個notify_*方法本質(zhì)上都是把對應的類型化事件投遞到SubscriptionBus見 context.py 中的_bus.publish(...)真正負責把事件送到每條流上的是 SDK 內(nèi)部的訂閱總線。MCPServer替你承擔了subscriptions/listen的服務器職責。協(xié)議層面的義務全部由 SDK 完成確認幀acknowledgment必須是流的第一幀按流做事件過濾每一幀上都攜帶訂閱標識符subscription id。線上的樣子確認幀與事件幀當過濾器指定了board://sprint的流在complete_task執(zhí)行后線上的幀如下{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}}}注意更新幀并不攜帶看板內(nèi)容本身。每一幀都在_meta中攜帶 listen 請求的 JSON-RPC id這個 id 就是訂閱標識符。它由客戶端鑄造Python 的Client使用listen-1這類字符串其他客戶端可能用整數(shù)。SDK 中SUBSCRIPTION_ID_META_KEY io.modelcontextprotocol/subscriptionId就是這段元數(shù)據(jù)的鍵名定義在 src/mcp/shared/subscriptions.py。過濾器就是契約只送達被請求的內(nèi)容過濾器是一份合同。一條請求了工具列表變化 一個資源 URI的流只會收到這兩種事件絕不會多。你發(fā)布一個提示詞變化這條流保持沉默。一個容易踩坑的細節(jié)MCPServer對資源 URI 按精確字符串匹配。因此訂閱了board://sprint的流聽不到關于board://sprint/tasks/1的變化。規(guī)范允許服務器報告訂閱 URI 的子資源的變化MCPServer從不這樣做但客戶端是按可能收到來構建的——所以客戶端讀取事件時應讀取event.uri而非假設具體是哪個資源動了詳見客戶端文檔 docs/client/subscriptions.md。底層過濾邏輯集中在 src/mcp/shared/subscriptions.py 的event_matchesToolsListChanged只在honored.tools_list_changed is True時通過ResourceUpdated只在該 URI 落在已確認的訂閱 URI 集合內(nèi)時通過——服務器投遞與客戶端接收共用同一個準入謂詞。這條流不是什么它不是重放日志replay log。斷掉的流就消失了斷連期間發(fā)布的事件不會被排隊??蛻舳酥匦?listen 并重新拉取數(shù)據(jù)。它不是 2025 時代的老路徑。調(diào)用過resources/subscribe的舊客戶端由ctx.session.send_resource_updated(uri)服務。notify_*方法只送達subscriptions/listen流context.py 的 docstring 明確說明了這一區(qū)別。誰能看用 middleware 把關訪問權限默認情況下每一個被請求的種類和 URI 都會被滿足任何調(diào)用者都可以觀看你發(fā)布的任何 URI。此時沒有人咨詢你的讀處理器——因為沒有人讀取。一個會被你的files://{name}處理器拒之門外的調(diào)用者仍然可以打開一條訂閱files://payroll.csv的流得知文件變了以及何時變的。它永遠不會得到內(nèi)容也無法探測哪些 URI 存在因為未知 URI 同樣被滿足只是永遠不會觸發(fā)。這個泄漏面很窄但真實存在——在從多租戶服務器發(fā)布按用戶區(qū)分的 URI 之前務必加上訪問檢查。這道檢查就是 middleware。它在 SDK 確認請求之前看到subscriptions/listen請求并在調(diào)用者請求了任何無權讀取的內(nèi)容時拒絕。完整示例見 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客戶端收到這個錯誤、得不到任何流而連接繼續(xù)存活。務必保持錯誤信息統(tǒng)一、不提及任何 URI這樣一次拒絕永遠不會向調(diào)用者證實哪些 URI 是受保護的。用同一個can_access(user, uri)回答兩個問題資源處理器在resources/read時調(diào)用它middleware 在subscriptions/listen時調(diào)用它。把示例里的表換成數(shù)據(jù)庫或你的 RBAC 系統(tǒng)兩個路徑始終保持一致。決定在流的整個生命周期內(nèi)生效沒有逐事件的復查。如果調(diào)用者的訪問權可能在流中途失效例如即將過期的 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(...)時會發(fā)出請求并等待服務器的確認因此當代碼塊開始時流已經(jīng)處于活躍狀態(tài)每一個類型化事件都是重新拉取數(shù)據(jù)的信號而不是攜帶數(shù)據(jù)的負載。這就是整個契約的完整展示??蛻舳藗鹊母鄡?nèi)容如何讓觀察者與主流程并行、流的結束與重新監(jiān)聽等在獨立的 docs/client/subscriptions.md 頁面中迭代產(chǎn)出四種類型化事件ToolsListChanged、PromptsListChanged、ResourcesListChanged、ResourceUpdated(uri...)重復的未消費事件會合并句柄提供sub.honored服務器確認的過濾器與sub.subscription_idlisten 請求的 id用于多路復用多條并發(fā)訂閱等屬性??邕M程擴展實現(xiàn)你自己的SubscriptionBus發(fā)布從你的 handler 到打開的流走的是一條SubscriptionBus。默認的總線在內(nèi)存中一個進程、該進程內(nèi)的所有流。在負載均衡器后面跑多副本之前這是完全正確的答案——因為一旦多副本客戶端的流被固定在某一個副本上而另一個副本上的發(fā)布必須能夠到達它。這個接縫由你來實現(xiàn)在你的 pub/sub 后端之上實現(xiàn)兩個方法。示例以 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是你的函數(shù)每個副本上負責解碼到達消息并調(diào)用每個已注冊 listener 的 reader 任務也是你的。Listener 是同步的、不得拋出異常、運行在服務器的事件循環(huán)上??偩€搬運的是類型化的ServerEvent值——四個小型 dataclass永遠不會是 JSON-RPC。打標stamping、過濾、流的生命周期都留在 SDK 里因此你的總線實現(xiàn)不可能破壞協(xié)議它只能把事件在進程間搬來搬去。SubscriptionBus在源碼中就是一個兩方法的Protocolpublish為異步、subscribe為同步本地注冊并返回冪等的 unsubscribe 可調(diào)用對象定義見 src/mcp/server/subscriptions.py其內(nèi)建實現(xiàn)InMemorySubscriptionBus在 同一文件 中發(fā)布時逐個調(diào)用 listener單個 listener 拋異常會被記錄并跳過隔離扇出邊界結束時還會做一次 checkpoint 讓事件流有排水機會。在請求之外發(fā)布持有總線引用要想在請求之外發(fā)布例如 lifespan 任務、webhook 等你需要自己構造總線以持有引用。MCPServer在你什么都不傳時會內(nèi)部構建一條但不會對外暴露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同樣的總線模式也出現(xiàn)在部署文檔中無論事件由哪臺服務器發(fā)布、流掛在哪個服務器對象上扇出本身不關心這些——同一進程內(nèi)的兩個MCPServer共享一條InMemorySubscriptionBus時行為已然如此在一個上開流、在另一個上發(fā)布流能聽到??缯鎸嵾M程時 SDK 不附帶任何可用的總線SubscriptionBus就是留給你在自己的 pub/sub 后端Redis、NATS 或任何你已經(jīng)在跑的上實現(xiàn)的接縫詳見 docs/run/deploy.md 的 Change notifications across replicas 一節(jié)。低層組合在沒有預接線的地方自己組裝在低層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, )三個要點總線歸你所有所以你直接向它發(fā)布await bus.publish(ResourceUpdated(uri...))。把它放在 handler 夠得著的地方——示例放在模塊級更大的應用放在 lifespan 里。ListenHandler(bus)就是MCPServer注冊的同一個 handleron_subscriptions_listen只是一個普通的 handler 槽位。在這個槽位放入你自己的可調(diào)用對象以獲得不同的語義屆時規(guī)范義務就轉移到你身上先確認、給每一幀打訂閱 id、絕不投遞過濾器之外的內(nèi)容。ListenHandler.close()優(yōu)雅地結束每一條打開的流每條流都會以 listen 請求的結果作為最后一幀收到它——這正是規(guī)范中服務器有意結束訂閱的表達方式。注意該方法在流沖刷完成之前就返回所以在拆除傳輸之前要給它們留一點時間。不調(diào)用它流會在客戶端斷開時結束。ListenHandler在 src/mcp/server/subscriptions.py 中的實現(xiàn)還透露出幾個實用細節(jié)單次調(diào)用即一條訂閱流構造參數(shù)max_subscriptions1024限制并發(fā)流數(shù)超限以INTERNAL_ERROR拒絕發(fā)生在確認幀之前max_buffered_events1024限制每條流的積壓事件數(shù)——當流的積壓到達上限時該流會被結束客戶端重新 listen 并重新拉取沒有重放所以不丟任何積壓之外的東西事件通過有界內(nèi)存流緩沖發(fā)布者永遠不會被慢消費者阻塞。訂閱發(fā)生在發(fā)送確認幀之前因此確認寫掛起期間發(fā)布的事件會被緩沖而不是丟失——確認幀仍是第一幀因為只有 handler 任務本身寫這條流??偨Y客戶端以一次subscriptions/listen請求加入響應即流服務它serving是內(nèi)置能力。你用ctx.notify_*發(fā)布SDK 完成打標、過濾和生命周期管理。事件是信號而非負載兩端都會重新拉取數(shù)據(jù)??蛻舳艘粋染褪莂sync with client.listen(...)詳見 docs/client/subscriptions.md。在低層Server上你自行組裝同樣的部件一條總線、ListenHandler(bus)、on_subscriptions_listen槽位。水平擴展 實現(xiàn)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 服務端訂閱機制實戰(zhàn)subscriptions/listen 事件流、過濾器與多進程擴展MCP Python SDK 服務端訂閱機制實戰(zhàn)subscriptions/listen 事件流、過濾器與多進程擴展 服務器目錄并非一成不變工具會在運行時出人工智能MCP 服務MCP ClientsMCP Python SDK 訂閱機制全解析從 subscriptions/listen 流式通知到多副本擴展MCP Python SDK 訂閱機制全解析從 subscriptions/listen 流式通知到多副本擴展 服務器的目錄并非一成不變工具會在運行時出現(xiàn)人工智能MCP 服務MCP Clients創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考