Go gorilla/websocket同端口RabbitMQ消费者发送WebSocket消息报错
完全不需要单独启动监听其他端口的Go程序,所有逻辑都可以在同一个main包、同一个进程内实现。你之前尝试用dial方式连自己的WebSocket服务发消息属于绕弯路,还踩了几个明显的代码错误,根本没必要这么做。
先把你之前碰到的几个错误根因说清楚:
dial:websocket: bad handshake 404:Gin路由注册的WebSocket路径和你dial时填写的Path不匹配,服务端找不到对应接口直接返回404,握手自然失败。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(¤tConn) }() 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, ¤tConn) connLock.Unlock() go handleIO(¤tConn, 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

