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

Go gorilla/websocket同端口RabbitMQ消费者发送WebSocket消息报错

问题结论

完全不需要单独启动监听其他端口的Go程序,所有逻辑都可以在同一个main包、同一个进程内实现。你之前尝试用dial方式连自己的WebSocket服务发消息属于绕弯路,还踩了几个明显的代码错误,根本没必要这么做。


你之前报错的原因

先把你之前碰到的几个错误根因说清楚:

  1. dial:websocket: bad handshake 404:Gin路由注册的WebSocket路径和你dial时填写的Path不匹配,服务端找不到对应接口直接返回404,握手自然失败。
  2. ERROR invalid character 'T' looking for beginning of value:你dial成功后直接发送了纯文本内容,但服务端handleIO逻辑里是用ReadJSON按JSON格式解析客户端消息的,纯文本不是合法JSON结构,解析第一个普通字符就会抛这个错。就算你要走dial方案,发送的内容也必须符合payload结构体的JSON格式,比如{"Message":"测试消息"}。

就算你把上面两个问题都修复,dial方案依然是下策:相当于你把Bot当成一个普通客户端连自己的服务,要额外维护长连接存活、断连重连逻辑,还多了一层不必要的本地网络转发开销,完全没必要。


正确实现方案

直接复用你现有的全局连接池,RabbitMQ消费者和Web服务在同一个进程内启动,消费到消息后直接调用广播逻辑推给所有在线用户即可,不需要走WebSocket连接。
注意你现有代码有两个严重的并发隐患,必须先修复,否则上线后会随机panic:

  • 全局connections切片会被多个WebSocket连接的读写协程、MQ消费协程同时增删改查,没有锁保护会出现竞态问题
  • gorilla/websocket的Conn不支持并发写入,你现在直接在广播逻辑里循环调用WriteJSON,多协程同时写同一个连接会直接触发panic

第一步:修复并发问题

先改造连接结构体,增加专用写通道和全局锁:

import "sync"

var (
    connections = make([]*webSocketConnection, 0)
    connLock    sync.RWMutex // 连接池读写锁
)

type webSocketConnection struct {
    *websocket.Conn
    Username string
    sendChan chan response // 每个连接单独的写消息通道,保证串行写入
}

在Execute函数初始化连接时,初始化写通道并启动专属写协程,保证所有写入操作都在同一个协程内执行:

func Execute(c *gin.Context, db repository.GormDB, qBroker *queue.Broker) {
    upgrader := websocket.Upgrader{
        ReadBufferSize:  maxMessageSize,
        WriteBufferSize: maxMessageSize,
        // 本地测试可以放开跨域,生产环境按需配置
        CheckOrigin: func(r *http.Request) bool { return true },
    }

    currentGorillaConn, err := upgrader.Upgrade(c.Writer, c.Request, nil)
    if err != nil {
        http.Error(c.Writer, "Could not open websocket connection", http.StatusBadRequest)
        return
    }
    username := c.Query("username")
    currentConn := webSocketConnection{
        Conn:     currentGorillaConn,
        Username: username,
        sendChan: make(chan response, 256), // 带缓冲,避免短时间消息堆积
    }

    // 启动专属写协程,所有发往该连接的消息都走这个协程,避免并发写
    go func() {
        defer func() {
            currentConn.Close()
            ejectConnection(&currentConn)
        }()
        for resp := range currentConn.sendChan {
            if err := currentConn.WriteJSON(resp); err != nil {
                log.Printf("write to user %s failed: %v", currentConn.Username, err)
                return
            }
        }
    }()

    connLock.Lock()
    connections = append(connections, &currentConn)
    connLock.Unlock()

    go handleIO(&currentConn, db, qBroker)
}

改造连接移除、广播函数,加锁保护,改为通过通道发送消息:

func ejectConnection(currentConn *webSocketConnection) {
    connLock.Lock()
    defer connLock.Unlock()
    // 关闭写通道,触发写协程退出
    close(currentConn.sendChan)
    filtered := gubrak.From(connections).Reject(func(each *webSocketConnection) bool {
        return each == currentConn
    }).Result()
    connections = filtered.([]*webSocketConnection)
}

func broadcastMessage(from string, msgType string, content string, exclude *webSocketConnection) {
    resp := response{
        From:    from,
        Type:    msgType,
        Message: content,
    }
    connLock.RLock()
    defer connLock.RUnlock()
    for _, eachConn := range connections {
        if eachConn == exclude {
            continue
        }
        // 非阻塞发送,通道满了直接丢弃消息并踢出异常连接
        select {
        case eachConn.sendChan <- resp:
        default:
            log.Printf("user %s send buffer full, close connection", eachConn.Username)
            // 异步处理断开,避免读锁重复加锁
            go ejectConnection(eachConn)
        }
    }
}

对应修改handleIO里的广播调用,顺便修复你原来/stock分支messageEntity未赋值的bug:

func handleIO(currentConn *webSocketConnection, db repository.GormDB, qBroker *queue.Broker) {
    defer func() {
        if r := recover(); r != nil {
            log.Println("ERROR", fmt.Sprintf("%v", r))
        }
        ejectConnection(currentConn)
    }()

    // 新用户加入广播
    broadcastMessage(currentConn.Username, messageNewUser, "", currentConn)
    joinMsg := message.NewMessage(currentConn.Username, messageNewUser, fmt.Sprintf("User %s: connected", currentConn.Username))
    db.Create(joinMsg)
    broadcastMessage("System", messageNewUser, joinMsg.Message, nil)

    for {
        payload := payload{}
        err := currentConn.ReadJSON(&payload)
        if err != nil {
            if strings.Contains(err.Error(), "websocket: close") {
                leaveTip := message.NewMessage(currentConn.Username, messageLeave, fmt.Sprintf("User %s: disconnect", currentConn.Username))
                db.Create(leaveTip)
                broadcastMessage("System", messageLeave, leaveTip.Message, nil)
                return
            }
            log.Println("ERROR", err.Error())
            continue
        }

        trimStr := strings.TrimSpace(payload.Message)
        splitStr := strings.Split(trimStr, "=")
        if splitStr[0] == "/stock" && len(splitStr)>=2 {
            _ = qBroker.PublishMessage("bot-send", splitStr[1])
            tipMsg := message.NewMessage(currentConn.Username, messageChat, fmt.Sprintf("正在查询%s的行情,请稍候...", splitStr[1]))
            db.Create(tipMsg)
            broadcastMessage(currentConn.Username, messageChat, tipMsg.Message, currentConn)
        } else {
            chatMsg := message.NewMessage(currentConn.Username, messageChat, payload.Message)
            db.Create(chatMsg)
            broadcastMessage(currentConn.Username, messageChat, chatMsg.Message, currentConn)
        }
    }
}

第二步:同进程启动RabbitMQ消费者

直接在main函数里启动Web服务前,开一个goroutine跑MQ消费逻辑即可,不需要额外端口、不需要单独进程:

func main() {
    // 你原来的初始化逻辑:加载配置、初始化DB、初始化RabbitMQ客户端、初始化Gin
    r := gin.Default()
    db := initDB() // 替换成你实际的DB初始化逻辑
    qBroker := queue.NewBroker() // 替换成你实际的MQ初始化逻辑

    // 注册WebSocket路由,注意路径要和前端请求地址一致,避免404
    r.GET("/socket/ws", func(c *gin.Context) {
        socket.Execute(c, db, qBroker)
    })

    // 启动Bot消费者协程,和Web服务同进程运行
    go func() {
        consumer, err := qBroker.Consume("bot-send") // 监听你之前发股票指令的队列
        if err != nil {
            log.Fatalf("start bot mq consumer failed: %v", err)
        }
        for mqMsg := range consumer {
            stockCode := string(mqMsg.Body)
            // 这里替换成你实际的股票查询、Bot回复生成逻辑
            botReply := fmt.Sprintf("股票%s当前价格为123.45元,涨跌幅+2.1%%", stockCode)
            // 消息存库
            botMsg := message.NewMessage("Bot", messageChat, botReply)
            db.Create(botMsg)
            // 广播给所有在线用户,exclude传nil表示不跳过任何连接
            broadcastMessage("Bot", messageChat, botReply, nil)
            _ = mqMsg.Ack(false)
        }
    }()

    // 启动Web服务
    log.Fatal(r.Run(":8081"))
}

补充说明

如果你后续确实需要让Bot支持跨服务部署,再考虑用dial方式作为普通客户端接入即可,届时注意发送的消息必须是符合payload结构的JSON格式,同时要给这个Bot客户端加心跳、断连重连逻辑,保证可用性。单实例部署场景下,直接操作连接池是性能最高、维护成本最低的方案。

内容的提问来源于stack exchange,提问作者Wincenty Bertoni Lech

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 18:03:15