第7讲:通知机制与实时性——MCP 的事件驱动能力

发布时间:2026/9/11 18:38:17

第7讲:通知机制与实时性——MCP 的事件驱动能力 一、为什么需要通知前六讲的交互模式全是Client 主动拉取PollingClient 调用tools/list获取工具列表Client 调用resources/read读取日志Client 调用tools/call执行动作这种模式在大多数场景下够用但有三个典型痛点场景Polling 的问题订单状态从“支付中”变为“已支付”Agent 不知道除非每隔几秒轮询一次运维工具新增了一个restart_podAgent 还在用旧的工具列表缓存直到缓存过期日志文件追加了新内容Agent 必须反复读取整个文件浪费带宽MCP 的通知机制notifications/*解决了这个问题Server 在状态变化时主动推送消息给 ClientClient 收到后决定是否刷新缓存或采取行动。二、MCP 通知的三类场景MCP 协议定义了三种标准通知通知方法触发时机Client 的响应notifications/tools/list_changedServer 的工具列表发生变化清除工具缓存下次tools/list重新拉取notifications/resources/list_changedServer 的资源列表或模板发生变化清除资源缓存notifications/initializedClient 初始化完成后通知 ServerServer 开始发送其他通知此外MCP 允许自定义通知notifications/*例如notifications/order/status_changed订单状态变更notifications/log/appended日志文件追加三、Go 实现带通知机制的 MCP Server下面扩展第4讲的日志 Resource Server增加工具列表变更通知模拟管理员新增工具日志文件追加通知模拟 tail -f通过 SSE 通道推送给 Clientpackage main import ( encoding/json fmt log net/http os path/filepath strings sync time ) // ---- MCP 消息结构 ---- type JSONRPCReq struct { JSONRPC string json:jsonrpc ID int json:id Method string json:method Params json.RawMessage json:params,omitempty } type JSONRPCResp struct { JSONRPC string json:jsonrpc ID int json:id Result interface{} json:result,omitempty Error *RPCError json:error,omitempty } type RPCError struct { Code int json:code Message string json:message } // ---- 通知消息JSON-RPC 通知格式无 id---- type Notification struct { JSONRPC string json:jsonrpc Method string json:method Params interface{} json:params,omitempty } // ---- 工具定义 ---- type ToolSpec struct { Name string json:name Description string json:description InputSchema any json:inputSchema } // ---- 资源定义 ---- type ResourceTemplate struct { URITemplate string json:uriTemplate Name string json:name Description string json:description MimeType string json:mimeType } // ---- SSE 客户端管理 ---- type SSEClientManager struct { mu sync.RWMutex clients map[chan Notification]bool } func NewSSEClientManager() *SSEClientManager { return SSEClientManager{ clients: make(map[chan Notification]bool), } } func (m *SSEClientManager) Subscribe() chan Notification { ch : make(chan Notification, 100) m.mu.Lock() m.clients[ch] true m.mu.Unlock() return ch } func (m *SSEClientManager) Unsubscribe(ch chan Notification) { m.mu.Lock() delete(m.clients, ch) m.mu.Unlock() close(ch) } func (m *SSEClientManager) Broadcast(notif Notification) { m.mu.RLock() defer m.mu.RUnlock() for ch : range m.clients { select { case ch - notif: default: // 客户端消费太慢丢弃消息 log.Println(通知队列满丢弃消息:, notif.Method) } } } // ---- MCP Server带通知能力---- type NotifyingMCPServer struct { sseManager *SSEClientManager tools map[string]ToolSpec mu sync.RWMutex } func NewNotifyingMCPServer(sseManager *SSEClientManager) *NotifyingMCPServer { return NotifyingMCPServer{ sseManager: sseManager, tools: map[string]ToolSpec{ current_time: { Name: current_time, Description: 返回当前 UTC 与北京时间, InputSchema: map[string]interface{}{ type: object, properties: map[string]interface{}{}, }, }, }, } } // 模拟管理员动态添加工具 func (s *NotifyingMCPServer) AddTool(name, desc string) { s.mu.Lock() s.tools[name] ToolSpec{ Name: name, Description: desc, InputSchema: map[string]interface{}{ type: object, properties: map[string]interface{}{}, }, } s.mu.Unlock() // 广播工具列表变更通知 s.sseManager.Broadcast(Notification{ JSONRPC: 2.0, Method: notifications/tools/list_changed, Params: map[string]interface{}{new_tool: name}, }) log.Printf(已添加工具 %s 并广播通知, name) } func (s *NotifyingMCPServer) Handle(req JSONRPCReq) JSONRPCResp { switch req.Method { case initialize: return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ protocolVersion: 2026-07-28, capabilities: map[string]interface{}{ tools: map[string]interface{}{}, resources: map[string]interface{}{}, }, }, } case tools/list: s.mu.RLock() specs : make([]ToolSpec, 0, len(s.tools)) for _, t : range s.tools { specs append(specs, t) } s.mu.RUnlock() return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{tools: specs}, } case tools/call: var params struct { Name string json:name Arguments json.RawMessage json:arguments } json.Unmarshal(req.Params, params) s.mu.RLock() _, exists : s.tools[params.Name] s.mu.RUnlock() if !exists { return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Error: RPCError{Code: -32601, Message: 未知工具}, } } return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ content: []map[string]interface{}{ {type: text, text: fmt.Sprintf(工具 %s 执行成功, params.Name)}, }, }, } case resources/list: return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ resources: []interface{}{}, resourceTemplates: []ResourceTemplate{ { URITemplate: logs://{namespace}/{pod}, Name: Pod 日志, Description: 按 namespace 和 pod 名称读取容器日志, MimeType: text/plain, }, }, }, } case resources/read: var params struct { URI string json:uri } json.Unmarshal(req.Params, params) // 简化实现返回模拟日志 if !strings.HasPrefix(params.URI, logs://) { return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Error: RPCError{Code: -32603, Message: 不支持的 URI}, } } return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ contents: []map[string]interface{}{ {uri: params.URI, mimeType: text/plain, text: fmt.Sprintf([%s] 模拟日志内容, time.Now().Format(time.RFC3339))}, }, }, } default: return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Error: RPCError{Code: -32601, Message: 不支持的方法}, } } } // ---- HTTP Transport双端点/mcp 用于 JSON-RPC/events 用于 SSE 通知---- func main() { sseManager : NewSSEClientManager() server : NewNotifyingMCPServer(sseManager) mux : http.NewServeMux() // JSON-RPC 端点 mux.HandleFunc(/mcp, func(w http.ResponseWriter, r *http.Request) { var req JSONRPCReq json.NewDecoder(r.Body).Decode(req) resp : server.Handle(req) w.Header().Set(Content-Type, application/json) json.NewEncoder(w).Encode(resp) }) // SSE 端点Client 订阅通知 mux.HandleFunc(/events, func(w http.ResponseWriter, r *http.Request) { flusher, ok : w.(http.Flusher) if !ok { http.Error(w, Streaming unsupported, http.StatusInternalServerError) return } w.Header().Set(Content-Type, text/event-stream) w.Header().Set(Cache-Control, no-cache) w.Header().Set(Connection, keep-alive) ch : sseManager.Subscribe() defer sseManager.Unsubscribe(ch) // 发送初始连接成功消息 fmt.Fprintf(w, event: connected\ndata: {\status\:\ok\}\n\n) flusher.Flush() ctx : r.Context() for { select { case -ctx.Done(): log.Println(SSE 客户端断开连接) return case notif : -ch: data, _ : json.Marshal(notif) fmt.Fprintf(w, event: message\ndata: %s\n\n, data) flusher.Flush() } } }) // 模拟每 15 秒自动添加一个新工具演示通知 go func() { toolNames : []string{ping_host, check_disk, restart_service, tail_log} idx : 0 for { time.Sleep(15 * time.Second) if idx len(toolNames) { server.AddTool(toolNames[idx], fmt.Sprintf(模拟工具 #%d, idx1)) idx } } }() // 模拟每 10 秒广播一条日志追加通知 go func() { for { time.Sleep(10 * time.Second) sseManager.Broadcast(Notification{ JSONRPC: 2.0, Method: notifications/log/appended, Params: map[string]interface{}{ uri: logs://production/web-abc, timestamp: time.Now().Format(time.RFC3339), line: fmt.Sprintf(新日志条目 at %s, time.Now().Format(time.RFC3339Nano)), }, }) } }() log.Println(MCP 通知 Server 启动:) log.Println( JSON-RPC: http://localhost:8083/mcp) log.Println( SSE: http://localhost:8083/events) log.Fatal(http.ListenAndServe(:8083, mux)) }四、Client 端对接通知Client 需要同时维护两个连接HTTP POST 到/mcp发送 JSON-RPC 请求SSE 连接到/events接收 Server 推送的通知// ---- MCP Client 带通知订阅 ---- type NotifyingMCPClient struct { serverURL string token string client *http.Client mu sync.RWMutex toolsCache []ToolSpec cacheExpiry time.Time } func (c *NotifyingMCPClient) StartEventLoop(ctx context.Context) { // 建立 SSE 连接 req, _ : http.NewRequestWithContext(ctx, GET, c.serverURL/events, nil) req.Header.Set(Authorization, Bearer c.token) resp, err : c.client.Do(req) if err ! nil { log.Printf(SSE 连接失败: %v, err) return } defer resp.Body.Close() reader : bufio.NewReader(resp.Body) for { line, err : reader.ReadString(\n) if err ! nil { log.Printf(SSE 读取失败: %v, err) return } if strings.HasPrefix(line, data: ) { var notif Notification json.Unmarshal([]byte(line[6:]), notif) switch notif.Method { case notifications/tools/list_changed: log.Println(收到工具列表变更通知清除缓存) c.mu.Lock() c.toolsCache nil c.cacheExpiry time.Time{} c.mu.Unlock() case notifications/log/appended: log.Printf(收到日志追加通知: %v, notif.Params) // 可以选择主动刷新相关 Resource } } } }五、通知 vs Polling 的选择场景推荐方式原因工具列表变更通知notifications/tools/list_changed变更频率极低通知即可无需轮询日志文件追加通知 增量读取通知告知有新内容Client 再用resources/read读增量订单状态变更通知实时性要求高轮询延迟大配置项变更通知变更频率低通知即可指标监控数据Polling定期拉取数据持续变化通知会过于频繁经验法则状态变化频率低于每分钟1次的用通知高于每分钟1次的用 Polling 本地缓存。六、安全分层L6本讲在 L1-L5 基础上补充通知机制的安全关注点L6通知安全SSE 端点同样需要认证Bearer Token不能开放给未经授权的 Client通知内容不应包含敏感信息Token、密钥、个人数据Server 应限制每个 Client 的通知队列大小防止慢 Client 导致 Server 内存泄漏通知速率限制每秒不超过 10 条通知防止恶意 Server 淹没 ClientClient 应验证通知的来源确认来自可信的 Server防止伪造通知触发缓存清空七、延伸阅读MCP Specification – Notifications标准通知方法的定义和格式SSE (Server-Sent Events) 规范W3CEventSource API 和事件流格式WebSocket vs SSE 对比为什么 MCP 选择 SSE 而非 WebSocket单向通知 vs 双向通信八、下一讲预告第8讲MCP 生态——Hub、网关与联邦单个 MCP Server 好写但生产环境面对的是几十上百个 Server谁来管理它们的注册与发现谁来统一鉴权与限流不同团队的 Server 如何互操作下一讲介绍 MCP Hub、MCP Gateway 的概念以及如何用 Go 实现一个轻量级的 MCP 网关统一管理多个下游 Server。开发之余的小工具推荐处理 Base64、JSON 格式化、JWT 解析、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top。所有计算在浏览器完成文件不上服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。
延伸阅读

更多相关文章

2026/9/11 18:38:17

定时任务学习

一.理论基础前置复习:什么是完全二叉树:除了最后一层外其他层都达到最大节点数,且最后一层节点都靠左排列1.小顶堆小顶堆(Min-Heap)是一种特殊的完全二叉树结构,核心原则只有一条:堆序性质:任意一个父节点的值,都小于或等于它的子节点的值。也…

2026/9/11 18:38:17

PMSM离线参数辨识:从仿真模型到最小二乘实现

简介:永磁同步电机(PMSM)离线辨识是获取定子电阻与电感等关键参数、支撑矢量控制等策略优化的基础工作。这套仿真模型面向电机控制方向的研究者与工程师,提供了从数据采集、模型搭建到参数估计与验证的完整闭环。压缩包共5个文件&…

2026/9/11 19:33:25

mybatis使用笔记、生命周期、拦截器等、类型转换

文章目录打印sql日志mybatis-config.xml方式application.yml里面配置配置类配置方式生命周期生命周期-StatementHandler生命周期-ParameterHandler生命周期-Executor生命周期-ResultSetHandlermybatis拦截器其他扫描方式select日志可以打出来,update和delete语句未打…

2026/9/11 19:33:25

DevOps与商业场景实战:从工具链到价值实现

1. 项目概述 "开发运维与商业场景实战"这个标题让我想起了多年前第一次参与企业级项目部署时的场景。当时作为新人,我完全无法理解为什么一个简单的功能上线需要经过那么多复杂流程。直到后来自己踩过各种坑,才真正明白DevOps和商业场景结合的…

2026/9/11 19:33:25

增强缓存错误报告:从模糊日志到精准排障的可观测性实践

1. 为什么缓存错误报告需要“增强”:一次排查事故的反思在数据库或分布式系统的日常运维里,缓存承担着给热数据加速的重任。但缓存不是保险箱,它也会出错:数据满了被逐出、条目意外失效、并发写入互相覆盖、序列化反序列化异常………

2026/9/11 19:33:25

AutoMapper在.NET中的高效对象映射实践指南

1. AutoMapper核心价值与应用场景AutoMapper作为.NET生态中广泛使用的对象映射工具,其核心价值在于消除应用程序中繁琐的对象转换代码。在实际项目中,我们经常遇到DTO(Data Transfer Object)与领域模型、视图模型之间的相互转换需…

2026/9/11 19:33:25

Java开发者转型AI:工程思维与高薪路径

1. 项目概述:Java开发者转型AI的破局之道"Java老鸟破局:转身AI,让多年积累成为高薪底气"这个标题精准捕捉了当前技术圈的一个普遍焦虑——传统Java开发者如何应对AI浪潮的冲击。作为一名在Java和AI交叉领域深耕多年的从业者&#x…

2026/9/11 19:28:24

背包问题-分支限界法求解

1. 分支限界法的基本概念, 与背包问题实例, 相关于(10.2)1.1, 组合优化问题的基础概念。组合优化问题,是关乎在有限的解空间范围之内, 去找到可以满足某种约束条件的最优解的问题。这些问题通常出现在资源分配、路径规划和背包问题等场景中。…

2026/9/10 16:39:38

超人会飞不算本事:系统稳定依赖清晰规则与边界设计

开头先不绕弯子。“#斯坦李吐槽dc 所以超人是无缘无故会飞的嘛哈哈哈哈哈哈哈锤哥真是技术人才啊!#雷神 #复联”这类调侃式短标题,第一波冲击力在于它把两个宇宙的角色塞进同一个吐槽箱里,但细想一下就能发现,它真正碰到的根本不是…

2026/9/10 11:16:38

超人VS蜘蛛侠:拆解超级IP的影响力与传播方法论

把“蜘蛛侠 vs 超人”放在 CSDN 上聊,可能很多人第一反应是走错片场了。但如果把这两个角色看成“两个持续运营了 80 多年的文化产品”,你会发现,这场比较本质上是两个不同 IP 策略的长期结果对比:超人赢在定义了整个超级英雄题材…

2026/9/9 16:31:09

基于CNN的调制信号识别:MATLAB实现时频图分类实战

简介:本资源是一套面向通信工程与信号处理方向学习者、研究者的深度学习实践方案,聚焦调制信号自动检测与识别这一典型无线通信任务,解决传统方法依赖人工特征、低信噪比下性能下降等痛点。压缩包共12个文件(10.73MB)&…

2026/9/10 12:32:02

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/10 15:19:50

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/10 15:49:53

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码