You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 05:36:25