Golang WebSocket服务器仅发Ping消息,自定义Demo消息无法推送的解决咨询
问题
我开发的Golang WebSocket服务器已成功建立连接,但无法发送自定义消息,仅能发送Ping消息。相关代码如下:
type WebsocketServer struct { conn *websocket.Conn ticker time.Ticker wg sync.WaitGroup stopChan chan interface{} } // SetupWebSocketRoutes sets up the routes for WebSocket connections. func SetupWebSocketRoutes(app *fiber.App, ctx appcontext.ChannelContext, hub *WebsocketHub) { app.Use("/ws", func(c *fiber.Ctx) error { if websocket.IsWebSocketUpgrade(c) { return c.Next() } return fiber.ErrUpgradeRequired }) app.Get("/ws", websocket.New(func(c *websocket.Conn) { log.Println("WebSocket connection established") id, betChan := hub.RegisterClient() defer hub.UnregisterClient(id) log.Printf("Client registered with ID: %v", id) timer := time.NewTicker(8 * time.Second) defer timer.Stop() for { select { case bet, ok := <-betChan: if !ok { log.Println("WebSocket connection closed due to channel closure") c.Close() return } log.Println("Received bet update from channel") msg, err := json.Marshal(bet) if err != nil { log.Printf("Error marshaling bet: %v", err) continue } err = c.WriteMessage(websocket.TextMessage, msg) if err != nil { log.Printf("Error writing message to WebSocket: %v", err) c.Close() return } case <-timer.C: log.Println("Sending ping message to WebSocket") err := c.WriteMessage(websocket.PingMessage, nil) if err != nil { log.Printf("Error writing ping message to WebSocket: %v", err) c.Close() return } } } })) } func PopulateBetUpdatesChannel(betUpdatesChan chan<- *models.Bet) { for { bet := &models.Bet{ // msg } betUpdatesChan <- bet } }
目前仅能发送Ping消息,自定义的Demo Bet消息无法推送至客户端。该Demo消息仅用于测试,后续将替换为真实服务数据。请问如何修改代码,实现服务器正常发送消息且客户端能正确接收?
解决方案
1. 修复测试消息生成逻辑
当前PopulateBetUpdatesChannel生成的是空Bet结构体,且无延迟死循环会导致通道阻塞。修改函数填充测试数据并添加发送间隔:
func PopulateBetUpdatesChannel(betUpdatesChan chan<- *models.Bet) { ticker := time.NewTicker(5 * time.Second) // 每5秒发送一条测试消息 defer ticker.Stop() for range ticker.C { // 填充测试用Bet数据 bet := &models.Bet{ ID: 123, UserID: "test_user_001", Amount: 100.50, Status: "pending", CreatedAt: time.Now(), } betUpdatesChan <- bet log.Println("Sent test bet update to broadcast channel") } }
2. 启动测试消息生成协程
在程序初始化阶段(如main函数),用goroutine启动该函数,传入hub的全局广播通道:
func main() { // 初始化fiber app、WebsocketHub等逻辑 app := fiber.New() hub := NewWebsocketHub() // 启动测试消息生成协程 go PopulateBetUpdatesChannel(hub.BroadcastChan) // 注册WebSocket路由 SetupWebSocketRoutes(app, nil, hub) // 启动服务 log.Fatal(app.Listen(":3000")) }
3. 完善WebsocketHub的广播逻辑
确保Hub能将全局通道的消息广播到所有客户端,示例实现:
type WebsocketHub struct { clients map[string]chan<- *models.Bet BroadcastChan chan *models.Bet mu sync.RWMutex } func NewWebsocketHub() *WebsocketHub { hub := &WebsocketHub{ clients: make(map[string]chan<- *models.Bet), BroadcastChan: make(chan *models.Bet), } // 启动广播协程 go hub.run() return hub } func (h *WebsocketHub) run() { for bet := range h.BroadcastChan { h.mu.RLock() // 遍历所有客户端通道发送消息 for _, clientChan := range h.clients { select { case clientChan <- bet: default: log.Println("Client channel blocked, skip message") } } h.mu.RUnlock() } } func (h *WebsocketHub) RegisterClient() (string, <-chan *models.Bet) { h.mu.Lock() defer h.mu.Unlock() clientID := uuid.NewString() // 需引入github.com/google/uuid生成唯一ID clientChan := make(chan *models.Bet, 10) // 带缓冲通道避免阻塞 h.clients[clientID] = clientChan return clientID, clientChan } func (h *WebsocketHub) UnregisterClient(clientID string) { h.mu.Lock() defer h.mu.Unlock() if clientChan, ok := h.clients[clientID]; ok { close(clientChan) delete(h.clients, clientID) } }
4. 验证客户端接收逻辑
确保客户端正确监听文本消息,示例JS客户端:
const ws = new WebSocket('ws://localhost:3000/ws'); ws.onmessage = function(event) { const bet = JSON.parse(event.data); console.log('Received bet update:', bet); }; ws.onopen = function() { console.log('WebSocket connected'); }; ws.onerror = function(error) { console.error('WebSocket error:', error); };
内容的提问来源于stack exchange,提问作者flame_heart_person
相关产品推荐
相关产品推荐

