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

资讯详情

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

WebSocket 分布式集群广播:基于 Redis Pub/Sub 与消息幂等

WebSocket 分布式集群广播:基于 Redis Pub/Sub 与消息幂等 WebSocket 分布式集群广播基于 Redis Pub/Sub 与消息幂等在现代协同式多智能体MAS工作空间或企业级实时 AI 监控大盘中一个任务通常由多个 Agent 在后台并发推演并需要将最新的推理进度、生成的图表与状态日志实时多播/广播Multicast / Broadcast给房间内的所有协作者与观察者终端。当服务端由单台物理机演进为包含数十个节点的Kubernetes 分布式 WebSocket 集群时传统的单机内存连接管理如单机map[userId]*websocket.Conn会瞬间引发致命的**“跨节点消息寻址与广播孤岛困境”**场景痛点用户 A 连接在网关节点Pod-1上而用户 B 连接在网关节点Pod-2上此时后台的Planner_Agent在Pod-3上生成了一份重要的分析结论并试图向当前会话房间广播Pod-3本地内存中根本没有用户 A 和用户 B 的 Socket 连接句柄导致跨节点广播消息无法送达如果采用简单的网络广播在弱网重试下又会引发**“消息重复推送与页面频繁刷屏闪烁”**。如何构建一套**“基于 Redis Pub/Sub 的分布式跨节点消息总线 本地连接局部路由 客户端全局消息 ID 幂等去重Message Idempotency”的高可用分布式 WebSocket 广播架构**一、WebSocket 分布式跨节点集群广播全景架构模型[ 后台执行 Agent (Pod-3) 产出协同新事件: {msg_id: M_001, room: R_888, text: 分析就绪} ] │ ▼ (发布至全局 Redis Pub/Sub 通道: room:R_888) ┌────────────────────────────────────────────────────────────────────────┐ │ Redis 分布式消息广播总线 (Pub/Sub Bus) │ └─────────────────────────────────┬──────────────────────────────────────┘ │ (全集群所有网关 Pod 监听广播流) ┌────────────────────────┴────────────────────────┐ ▼ ▼ ┌─────────────────────────┐ ┌─────────────────────────┐ │ WebSocket 网关 Pod-1 │ │ WebSocket 网关 Pod-2 │ │ 本地内存查找: 发现 [用户A]│ │ 本地内存查找: 发现 [用户B]│ │ 动作: 向用户A Socket推流│ │ 动作: 向用户B Socket推流│ └────────────┬────────────┘ └────────────┬────────────┘ │ (通过 WS 下发) │ (通过 WS 下发) ▼ ▼ [ 终端用户 A 浏览器 ] [ 终端用户 B 浏览器 ] (前端基于 msg_id 做幂等去重0 重复渲染!)二、生产级 Go 语言分布式 WebSocket 广播器实现实操package broadcast import ( context encoding/json fmt sync github.com/gorilla/websocket github.com/redis/go-redis/v9 ) type BroadcastMessage struct { MessageID string json:msg_id // 全局唯一消息 ID用于前端幂等去重 RoomID string json:room_id Payload string json:payload } type DistributedWSGatewayNode struct { rdb *redis.Client nodeName string mu sync.RWMutex // 本地房间连接表: room_id - map[conn_ptr]bool localRoomConns map[string]map[*websocket.Conn]bool } func NewWSGatewayNode(rdbClient *redis.Client, nodeName string) *DistributedWSGatewayNode { node : DistributedWSGatewayNode{ rdb: rdbClient, nodeName: nodeName, localRoomConns: make(map[string]map[*websocket.Conn]bool), } // 启动后台监听 Redis 广播 go node.subscribeGlobalBroadcastEvents() return node } // RegisterLocalConn 将客户端 WebSocket 注册到本地房间 func (n *DistributedWSGatewayNode) RegisterLocalConn(roomID string, conn *websocket.Conn) { n.mu.Lock() defer n.mu.Unlock() if _, ok : n.localRoomConns[roomID]; !ok { n.localRoomConns[roomID] make(map[*websocket.Conn]bool) } n.localRoomConns[roomID][conn] true fmt.Printf( [%s] 客户端连接成功注册到房间: [%s]\n, n.nodeName, roomID) } // PublishEventToRoom 任何节点均可调用向指定房间发布广播事件 func (n *DistributedWSGatewayNode) PublishEventToRoom(ctx context.Context, msg BroadcastMessage) error { bytes, _ : json.Marshal(msg) channelName : fmt.Sprintf(ws_broadcast_room_%s, msg.RoomID) // 推入 Redis Pub/Sub由 Redis 负责通知全网所有网关 Pod return n.rdb.Publish(ctx, channelName, bytes).Err() } func (n *DistributedWSGatewayNode) subscribeGlobalBroadcastEvents() { ctx : context.Background() // 订阅所有房间通道 pubsub : n.rdb.PSubscribe(ctx, ws_broadcast_room_*) defer pubsub.Close() ch : pubsub.Channel() for msg : range ch { var broadcastMsg BroadcastMessage if err : json.Unmarshal([]byte(msg.Payload), broadcastMsg); err ! nil { continue } // 检查本地节点是否有属于该房间的活跃长连接 n.mu.RLock() conns, exists : n.localRoomConns[broadcastMsg.RoomID] if exists len(conns) 0 { // 本地并发向各个客户端 Socket 广播推送 for c : range conns { _ c.WriteMessage(websocket.TextMessage, []byte(broadcastMsg.Payload)) } fmt.Printf( [%s] 成功将事件 [%s] 本地推流给房间 [%s] 的 %d 个客户端\n, n.nodeName, broadcastMsg.MessageID, broadcastMsg.RoomID, len(conns)) } n.mu.RUnlock() } }三、前端基于全局msg_id的幂等去重渲染实现TypeScript在前端浏览器中利用Set或 LRU 缓存过滤重复网络重放帧const processedMsgIds new Setstring(); function handleIncomingWebSocketFrame(rawJson: string) { const event JSON.parse(rawJson); // 【核心幂等防线】若该消息已渲染过瞬间丢弃0 界面重绘抖动 if (processedMsgIds.has(event.msg_id)) { console.warn([幂等过滤] 丢弃重复到达的消息帧: ${event.msg_id}); return; } processedMsgIds.add(event.msg_id); // 执行真实 UI 渲染更新 renderMessageToChatUI(event.payload); }四、生产治理收益通过在多智能体协同底座中构建分布式 WebSocket 广播架构实现了全集群任意 Pod 之间毫秒级 5ms跨节点实时消息多播前端页面 100% 免疫了由于网络抖动重连引发的重复消息刷屏与闪烁支持 WebSocket 网关节点在 Kubernetes 中进行无状态平滑水平弹性扩缩容。
返回列表