Golang gorilla/websocket多客户端广播消息轮流接收问题修复
问题根因
消息轮流分发的核心原因是每个WebSocket连接的独立处理协程,都在直接调用handler.processor.TakeList()消费消息队列。TakeList()的逻辑是从队列取出消息后就将消息从队列移除,多个连接协程并发抢占队列消息时,自然会出现不同客户端各自抢到一部分消息的情况——当前实现本质是多消费者抢队列的负载均衡模式,根本不是广播模式。客户端代码逻辑无问题,不需要修改。
修正方案
广播模式的核心逻辑是:消息仅从队列消费一次,服务端维护所有在线客户端的连接集合,拿到新消息后遍历全部连接,给每个连接推送相同的消息副本,而非让每个连接自行到队列抢消息。
1. 实现线程安全的在线连接池
Go原生map非线程安全,搭配读写锁维护所有在线连接:
import "sync" var ( onlineClients = make(map[*websocket.Conn]struct{}) clientLock sync.RWMutex )
新客户端连接升级成功后加入连接池,连接断开时及时从池中移除:
conn, err := upgrader.Upgrade(w, r, nil) if err != nil { log.Print("WebSocket upgrade failed:", err) return } // 新连接入池 clientLock.Lock() onlineClients[conn] = struct{}{} clientLock.Unlock() // 连接关闭时出池 defer func() { clientLock.Lock() delete(onlineClients, conn) conn.Close() clientLock.Unlock() }()
2. 启动独立的广播消费协程
不要在单个连接的处理逻辑中消费消息队列,服务启动时单独启动一个后台协程,专门负责从队列取消息、遍历所有在线连接推送:
// 服务初始化阶段启动该协程 go func() { for { msgList, err := handler.processor.TakeList() if err != nil { log.Println("fetch broadcast msg failed:", err) time.Sleep(time.Millisecond * 100) // 避免出错时空转占CPU continue } // 遍历所有在线连接推送消息 clientLock.RLock() for conn := range onlineClients { if err := conn.WriteJSON(msgList); err != nil { log.Println("push msg to client failed:", err) // 推送失败先释放读锁,移除失效连接 clientLock.RUnlock() clientLock.Lock() delete(onlineClients, conn) _ = conn.Close() clientLock.Unlock() clientLock.RLock() } } clientLock.RUnlock() } }()
3. 清理单连接逻辑中的冗余代码
删除原有每个连接处理循环里调用TakeList()、WriteJSON的逻辑,单连接的处理逻辑只需要保留读取客户端上行消息、检测连接状态的部分即可。
注意事项
- WebSocket连接不支持并发写,上述实现中只有广播协程负责向连接写消息,不存在并发写风险;如果后续有其他场景需要向连接写数据,要为每个连接配置独立的写锁保证安全。
- 及时清理断连的客户端,避免连接池内存泄漏。
- 高并发场景下建议给每个连接配置带缓冲的发送队列,避免单个慢客户端阻塞整个广播流程。
内容的提问来源于stack exchange,提问作者TheQQWEETT
相关产品推荐
相关产品推荐

