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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 23:40:58