WebSocket 实时通信实现:群聊消息推送全链路解析

摘要: HTTP 是无状态的,但聊天室不能靠轮询撑着。本文以 CampusHub 校园活动平台的群聊服务为例,从 WebSocket 升级握手、Hub + Client 双层连接管理、JWT 鉴权、群聊广播,到异步消息保存队列(Worker Pool + 死信队列)、Redis Stream 消息中间件解耦、在线状态管理,完整拆解一个生产级 WebSocket 服务的实现路径。所有代码均来自真实项目。

标签: WebSocket, go-zero, 微服务, 实时通信, 群聊, Go语言, 实战

分类: 后端开发


引言

上一篇文章中,我们解决了微服务之间"怎么说话"的问题。但校园活动平台还有另一个核心需求:群聊

用户加入活动群聊后,能实时收到其他成员的消息——这是 HTTP 请求-响应模式天然无法实现的。让我们看一下常见的选择:

方案原理延迟服务器压力适用场景
短轮询客户端每 N 秒请求一次高(N 秒延迟)大量无效请求不推荐
长轮询请求挂起直到有数据连接占用多简单推送
SSE服务端单向推流较低仅需服务端推
WebSocket全双工长连接极低低(连接复用)双向实时通信

WebSocket 一次握手,之后全双工通信,是群聊场景的不二之选。但在微服务架构中,WebSocket 有三个额外挑战:

  1. 鉴权:WebSocket 握手是 HTTP 请求,但不能用普通的 JWT 中间件——连接建立后才能鉴权
  2. 多实例广播:消息发到实例 A,接收方的连接在实例 B,怎么跨实例推送?
  3. 消息可靠性:连接瞬间断了,消息还要不要保存到数据库?

CampusHub 的解决思路:同进程 Hub 管理连接 + Redis Stream 跨实例广播 + 异步队列保证消息落库

WebSocket

WebSocket

WebSocket

发布消息

订阅推送

异步写入

RPC 调用

状态存储

客户端 A

WebSocket 服务

客户端 B

客户端 C

Redis Stream
chat.message.new

SaveQueue
消息保存队列

Chat RPC
持久化到 MySQL

Redis
user:status:*

通过本文,你将学会

  • ✅ 用 gorilla/websocket 实现 HTTP → WebSocket 协议升级
  • ✅ 设计 Hub + Client 双层连接管理架构(并发安全)
  • ✅ 在 WebSocket 连接建立后完成 JWT 鉴权
  • ✅ 实现群聊消息广播(BroadcastToGroup)
  • ✅ 用 Worker Pool + 死信队列构建异步消息保存队列
  • ✅ 通过 Redis Stream 解耦实时推送与消息持久化
  • ✅ 用 Redis HMSet 管理用户在线状态
  • ✅ 实现心跳保活和优雅关闭

问题分析

群聊服务表面是"收发消息",实际要解决五个核心问题:

问题说明
🔌 连接管理同时在线几千个 WebSocket 连接,怎么组织?群聊成员怎么维护?
🔐 鉴权时机WebSocket 握手是 GET 请求,Token 放哪里?连接建立后怎么验证身份?
📢 消息广播一条消息要推给群里所有在线成员,怎么实现?多实例部署时怎么跨节点推?
💾 消息可靠性发消息和落库是同步还是异步?同步会增加延迟,异步怎么保证不丢?
💔 连接保活客户端网络波动、NAT 超时,长连接怎么保持?断开了怎么感知?

常见错误做法

连接用全局 map 直接访问 → 并发写 map 导致 panic,线上出过的经典事故
WebSocket 握手时在 URL 里带 Token → URL 会被日志记录,Token 泄露风险
发消息同步等待数据库写入 → 链路延迟从 1ms 飙到 50ms,用户体验直线下降
无心跳机制 → 客户端悄悄断线,服务端仍认为在线,消息推送到"黑洞"

正确做法

Hub 统一管理连接 → 所有连接操作通过 channel 序列化,无锁竞争
连接建立后发 auth 消息完成鉴权 → Token 在消息体里,不暴露在 URL
ACK 即时 + 异步落库 → 消息入队立即返回 ACK,后台 Worker 异步保存
Ping/Pong + 读超时 → 服务端定时 Ping,60 秒无响应主动断开


整体架构

CampusHub 的 WebSocket 服务支持两种部署模式:

app/chat/ws/         ← WebSocket 独立服务(本文重点)
app/chat/api/        ← chat-api(也可集成 WebSocket,按需选择)

独立部署的好处:WebSocket 连接是长连接,水平扩展时需要独立调度,和 HTTP 短连接服务混在一起不易运维。

服务内部的分层结构:

websocket.go         ← 入口:启动 HTTP 服务器,注册路由
├── handler/
│   └── ws_handler.go      ← 升级 HTTP → WebSocket,创建 Client
├── hub/
│   ├── hub.go             ← Hub:连接注册/注销,消息路由,群聊管理
│   └── client.go          ← Client:读写 pump,心跳,消息分发
├── logic/
│   └── message_logic.go   ← 业务逻辑:鉴权、消息处理、用户名查询
├── types/
│   └── types.go           ← 消息协议定义(WSMessage 统一格式)
├── queue/
│   └── save_queue.go      ← 异步消息保存队列(Worker Pool + 死信)
└── svc/
    └── service_context.go ← 服务上下文(RPC客户端、Redis、队列等)

完整实现

📌 示例 1(基础):协议升级 + Client 双 goroutine 模型

消息协议设计

所有 WebSocket 消息共用一个统一格式,通过 type 字段区分消息类型:

// 📁 app/chat/ws/internal/types/types.go

// MessageType 消息类型(客户端 ↔ 服务端)
type MessageType string

const (
    // 客户端 → 服务端
    TypePing        MessageType = "ping"         // 心跳
    TypeAuth        MessageType = "auth"         // 认证(连接后第一条消息)
    TypeSendMessage MessageType = "send_message" // 发送消息

    // 服务端 → 客户端
    TypePong           MessageType = "pong"            // 心跳响应
    TypeAuthSuccess    MessageType = "auth_success"    // 认证成功
    TypeAuthFailed     MessageType = "auth_failed"     // 认证失败
    TypeNewMessage     MessageType = "new_message"     // 新消息(群聊广播)
    TypeVerifyProgress MessageType = "verify_progress" // 认证进度推送
    TypeError          MessageType = "error"           // 错误消息
    TypeAck            MessageType = "ack"             // 消息确认
)

// WSMessage WebSocket 消息统一格式
type WSMessage struct {
    Type      MessageType     `json:"type"`           // 消息类型
    MessageID string          `json:"message_id"`     // 消息ID(用于去重和确认)
    Timestamp int64           `json:"timestamp"`      // Unix 时间戳
    Data      json.RawMessage `json:"data,omitempty"` // 消息体(不同类型结构不同)
}

// SendMessageData 客户端发送消息时的 Data 内容
type SendMessageData struct {
    GroupID  string `json:"group_id"`            // 目标群聊
    MsgType  int32  `json:"msg_type"`            // 1-文字 2-图片
    Content  string `json:"content,omitempty"`
    ImageURL string `json:"image_url,omitempty"`
}

💡 设计要点Data 字段使用 json.RawMessage 而非 interface{},避免二次序列化,减少内存分配。调用方直接将 JSON 字节写入,性能更好。

HTTP 升级为 WebSocket
// 📁 app/chat/ws/internal/handler/ws_handler.go

var upgrader = websocket.Upgrader{
    ReadBufferSize:  1024,
    WriteBufferSize: 1024,
    // 允许跨域(生产环境应限制 Origin)
    CheckOrigin: func(r *http.Request) bool {
        return true
    },
}

// WebSocketHandler 处理 WebSocket 连接升级
func WebSocketHandler(svcCtx *svc.ServiceContext, h *hub.Hub) http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        // 1. HTTP → WebSocket 协议升级(三次握手)
        conn, err := upgrader.Upgrade(w, r, nil)
        if err != nil {
            logx.Errorf("升级连接失败: %v", err)
            return // 升级失败,upgrader 会自动返回 400/403
        }

        // 2. 创建 Client(此时还未认证,无 userID)
        client := hub.NewClient(h, conn)

        // 3. 注册到 Hub(放入注册 channel,Hub 的 Run 循环处理)
        h.Register() <- client

        // 4. 启动读写双 goroutine
        go client.WritePump() // 负责从 send channel 取数据写入连接
        go client.ReadPump()  // 负责从连接读数据并分发处理

        logx.Info("新的 WebSocket 连接已建立")
    }
}

💡 为什么是双 goroutine? gorilla/websocket 的连接不是并发安全的——读和写必须在独立的 goroutine 中。ReadPump 阻塞等待客户端消息,WritePump 阻塞等待 send channel 中的数据,两者互不干扰。

Client:读写 pump + 心跳保活
// 📁 app/chat/ws/hub/client.go

const (
    writeWait      = 10 * time.Second        // 写操作超时
    pongWait       = 60 * time.Second        // 等待 Pong 的超时(60s 无响应则断开)
    pingPeriod     = (pongWait * 9) / 10     // Ping 间隔(54s,必须小于 pongWait)
    maxMessageSize = 512 * 1024              // 最大消息体积 512KB
)

// Client WebSocket 客户端(一个连接对应一个 Client)
type Client struct {
    hub      *Hub
    conn     *websocket.Conn
    send     chan []byte   // 待发送消息的缓冲 channel(容量 256)
    userID   string        // 认证后才有值
    groups   map[string]bool // 已加入的群聊
    mu       sync.RWMutex
    isAuthed bool
}

// ReadPump 读 goroutine:从连接读数据 → 分发给 Hub 处理
func (c *Client) ReadPump() {
    defer func() {
        c.hub.unregister <- c // 离开时触发注销
        c.conn.Close()
    }()

    c.conn.SetReadLimit(maxMessageSize)
    c.conn.SetReadDeadline(time.Now().Add(pongWait))

    // 收到 Pong 时重置读超时(保活机制的关键)
    c.conn.SetPongHandler(func(string) error {
        c.conn.SetReadDeadline(time.Now().Add(pongWait))
        return nil
    })

    for {
        _, message, err := c.conn.ReadMessage()
        if err != nil {
            if websocket.IsUnexpectedCloseError(err,
                websocket.CloseGoingAway, websocket.CloseAbnormalClosure) {
                logx.Errorf("WebSocket 异常关闭: %v", err)
            }
            break // 退出循环,触发 defer 中的注销
        }
        c.handleMessage(message)
    }
}

// WritePump 写 goroutine:从 send channel 取数据 → 写入连接
func (c *Client) WritePump() {
    ticker := time.NewTicker(pingPeriod) // 定时 Ping
    defer func() {
        ticker.Stop()
        c.conn.Close()
    }()

    for {
        select {
        case message, ok := <-c.send:
            c.conn.SetWriteDeadline(time.Now().Add(writeWait))
            if !ok {
                // send channel 已关闭(Hub 注销了这个 Client)
                c.conn.WriteMessage(websocket.CloseMessage, []byte{})
                return
            }
            w, err := c.conn.NextWriter(websocket.TextMessage)
            if err != nil {
                return
            }
            w.Write(message)

            // 批量写入:把队列里已积压的消息一并写出,减少系统调用次数
            n := len(c.send)
            for i := 0; i < n; i++ {
                w.Write([]byte{'\n'})
                w.Write(<-c.send)
            }
            w.Close()

        case <-ticker.C:
            // 定时发 Ping,触发客户端回 Pong,ReadPump 重置超时
            c.conn.SetWriteDeadline(time.Now().Add(writeWait))
            if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
                return // 写 Ping 失败,说明连接已断,退出
            }
        }
    }
}

心跳保活的时序:

服务端 ReadPump 客户端 服务端 WritePump 服务端 ReadPump 客户端 服务端 WritePump loop [每 54 秒] 60 秒未收到 Pong Ping(WebSocket 控制帧) Pong(自动回复) 重置读超时为 +60s 读超时,返回 error defer:触发 unregister

📌 示例 2(完整):Hub 连接管理 + JWT 鉴权 + 群聊广播

Hub:连接管理中心

Hub 是整个 WebSocket 服务的"中枢神经",所有连接的注册、注销、消息路由都经过它:

// 📁 app/chat/ws/hub/hub.go

// Hub 连接管理中心
type Hub struct {
    clients  map[string]*Client       // userID → Client(已认证的连接)
    groups   map[string]map[*Client]bool // groupID → 成员集合

    register   chan *Client
    unregister chan *Client

    messageHandler  MessageHandler   // 业务逻辑处理器(可替换,方便测试)
    messagingClient *messaging.Client // Redis Stream 消息中间件
    redisClient     *redis.Client     // 在线状态存储

    mu sync.RWMutex // 保护 clients 和 groups 的并发读写
}

// Run Hub 主循环(单 goroutine 处理注册/注销,避免 map 并发竞争)
func (h *Hub) Run(ctx context.Context) {
    go h.subscribeMessages(ctx) // 独立 goroutine 订阅消息中间件

    for {
        select {
        case client := <-h.register:
            h.registerClient(client)

        case client := <-h.unregister:
            h.unregisterClient(client)

        case <-ctx.Done():
            logx.Info("Hub 正在关闭")
            return
        }
    }
}

// registerClient 注册客户端(在 Run goroutine 中调用,无需额外加锁)
func (h *Hub) registerClient(client *Client) {
    h.mu.Lock()
    defer h.mu.Unlock()

    if client.userID != "" {
        // 同一用户重复连接,关闭旧连接(多端互踢,按需配置)
        if oldClient, exists := h.clients[client.userID]; exists {
            close(oldClient.send)
        }
        h.clients[client.userID] = client
        h.updateUserStatus(client.userID, true) // 更新在线状态到 Redis
        logx.Infof("用户 %s 上线", client.userID)
    }
}

// unregisterClient 注销客户端,清理所有群聊订阅
func (h *Hub) unregisterClient(client *Client) {
    h.mu.Lock()
    defer h.mu.Unlock()

    if _, exists := h.clients[client.userID]; exists {
        delete(h.clients, client.userID)
        close(client.send)

        // 从用户已加入的所有群聊中移除
        for groupID := range client.groups {
            if clients, ok := h.groups[groupID]; ok {
                delete(clients, client)
                if len(clients) == 0 {
                    delete(h.groups, groupID)
                }
            }
        }

        h.updateUserStatus(client.userID, false) // 更新离线状态
        logx.Infof("用户 %s 离线", client.userID)
    }
}

💡 并发安全的关键registerunregister 是 channel,修改 clients map 的操作全部在 Run 这一个 goroutine 里执行,无竞争。读操作(广播时遍历 clients)用 RWMutex 保护,允许多个 goroutine 并发读。

WebSocket 鉴权:连接后发 auth 消息

WebSocket 握手是 HTTP GET 请求,Token 放 URL 有泄露风险,放 Header 在浏览器端不易设置。CampusHub 的做法是:连接建立后,客户端主动发一条 auth 消息完成身份验证

// 📁 app/chat/ws/hub/client.go

// handleMessage 接收到消息后的处理逻辑
func (c *Client) handleMessage(message []byte) {
    var msg types.WSMessage
    if err := json.Unmarshal(message, &msg); err != nil {
        c.SendError(400, "消息格式错误")
        return
    }

    // 未认证状态下,只允许 ping 和 auth 两种消息
    if !c.isAuthed && msg.Type != types.TypeAuth && msg.Type != types.TypePing {
        c.SendError(401, "未授权,请先发送 auth 消息")
        return
    }

    c.hub.handleClientMessage(c, &msg)
}

// 📁 app/chat/ws/internal/logic/message_logic.go

// HandleAuth 处理认证消息
func (l *MessageLogic) HandleAuth(client *hub.Client, msg *types.WSMessage) error {
    var authData types.AuthData
    if err := json.Unmarshal(msg.Data, &authData); err != nil {
        return err
    }

    // 1. 解析 JWT Token,提取 userID
    userID, err := l.svcCtx.JwtAuth.ParseToken(authData.Token)
    if err != nil {
        client.SendMessage(&types.WSMessage{
            Type:      types.TypeAuthFailed,
            Timestamp: time.Now().Unix(),
            Data:      json.RawMessage(`{"message":"Token 无效或已过期"}`),
        })
        return err
    }

    // 2. 设置客户端身份,标记为已认证
    client.SetUserID(userID)
    client.SetAuthed(true)

    // 3. 发送认证成功响应
    successData, _ := json.Marshal(map[string]string{"user_id": userID})
    client.SendMessage(&types.WSMessage{
        Type:      types.TypeAuthSuccess,
        Timestamp: time.Now().Unix(),
        Data:      successData,
    })

    // 4. 异步加入用户的所有群聊(避免阻塞认证响应)
    go l.autoJoinUserGroups(client, userID)

    logx.Infof("用户 %s 认证成功", userID)
    return nil
}

认证完成后,autoJoinUserGroups 会调用 Chat RPC 获取用户的群聊列表,并将 Client 注册到对应的群聊分组中,后续广播时直接定向推送:

func (l *MessageLogic) autoJoinUserGroups(client *hub.Client, userID string) {
    userIDUint, _ := strconv.ParseUint(userID, 10, 64)

    // 调用 Chat RPC 获取用户加入的所有群聊
    resp, err := l.svcCtx.ChatRpc.GetUserGroups(l.ctx, &chat.GetUserGroupsReq{
        UserId:   userIDUint,
        Page:     1,
        PageSize: 100, // 最多订阅 100 个群聊
    })
    if err != nil {
        logx.Errorf("获取用户群聊列表失败: %v", err)
        return
    }

    for _, group := range resp.Groups {
        client.GetHub().AddClientToGroup(client, group.GroupId)
    }
    logx.Infof("用户 %s 自动订阅了 %d 个群聊", userID, len(resp.Groups))
}

完整的连接建立时序:

Chat RPC MessageLogic Hub Handler 客户端 Chat RPC MessageLogic Hub Handler 客户端 连接建立,启动 ReadPump/WritePump 认证完成,开始收发消息 GET /ws(HTTP 升级请求) upgrader.Upgrade() register <- client(未认证) 101 Switching Protocols {"type":"auth","data":{"token":"eyJ..."}} HandleAuth(client, msg) ParseToken → userID="12345" {"type":"auth_success","data":{"user_id":"12345"}} GetUserGroups(userId=12345) [groupA, groupB, groupC] AddClientToGroup(client, groupA/B/C)
群聊消息广播
// 📁 app/chat/ws/hub/hub.go

// BroadcastToGroup 向群聊内所有在线成员广播消息
func (h *Hub) BroadcastToGroup(groupID string, msg *types.WSMessage) {
    data, err := json.Marshal(msg)
    if err != nil {
        logx.Errorf("序列化广播消息失败: %v", err)
        return
    }

    h.mu.RLock() // 读锁:允许并发广播
    defer h.mu.RUnlock()

    clients, ok := h.groups[groupID]
    if !ok || len(clients) == 0 {
        return // 群聊无在线成员,跳过
    }

    for client := range clients {
        client.sendRaw(data) // 非阻塞投递到 send channel
    }
}

// SendToUser 向指定用户单播消息
func (h *Hub) SendToUser(userID string, msg *types.WSMessage) error {
    h.mu.RLock()
    client, ok := h.clients[userID]
    h.mu.RUnlock()

    if !ok {
        return ErrUserNotOnline
    }
    return client.SendMessage(msg)
}
在线状态管理

用户连接和断开时,Hub 自动更新 Redis 中的在线状态:

// 📁 app/chat/ws/hub/hub.go

// updateUserStatus 更新用户在线状态到 Redis(HMSet 存储多个字段)
func (h *Hub) updateUserStatus(userID string, isOnline bool) {
    ctx := context.Background()
    key := fmt.Sprintf("user:status:%s", userID)
    now := time.Now().Unix()

    data := map[string]interface{}{
        "is_online": isOnline,
        "last_seen": now,
    }
    if isOnline {
        data["last_online_at"] = now
    } else {
        data["last_offline_at"] = now
    }

    if err := h.redisClient.HMSet(ctx, key, data).Err(); err != nil {
        logx.Errorf("更新用户状态失败: %v", err)
        return
    }
    // 30 天自动过期,避免 Redis 积累过多僵尸数据
    h.redisClient.Expire(ctx, key, 30*24*time.Hour)
}

Redis 中存储的数据结构:

HGETALL user:status:12345
1) "is_online"
2) "true"
3) "last_seen"
4) "1740614400"
5) "last_online_at"
6) "1740614400"

📌 示例 3(优化):异步消息队列 + 消息中间件解耦

核心设计:ACK 先行,异步落库

用户发送一条消息时,最朴素的做法是:先写数据库,再广播。问题是数据库写入耗时 20-50ms,用户等待 ACK 的时间就是这么长。

CampusHub 的优化方案:立即返回 ACK,写库和广播完全异步

MySQL Worker Redis Stream SaveQueue MessageLogic 客户端 MySQL Worker Redis Stream SaveQueue MessageLogic 客户端 par [并行执行] send_message(groupId, content) 生成 messageID(UUID) ACK(延迟 ~1ms,用户无感知) Publish("chat.message.new", msgData) Push(SaveMessageTask) 订阅回调:BroadcastToGroup new_message(推送给群内所有成员) 分发给空闲 Worker Chat RPC → SaveMessage
发送消息逻辑(异步版本)
// 📁 app/chat/ws/internal/logic/message_logic.go

// HandleSendMessage 处理发送消息(异步落库版本)
func (l *MessageLogic) HandleSendMessage(client *hub.Client, msg *types.WSMessage) error {
    var sendData types.SendMessageData
    if err := json.Unmarshal(msg.Data, &sendData); err != nil {
        return err
    }

    messageID := uuid.New().String()
    now := time.Now().Unix()
    senderID, _ := strconv.ParseUint(client.GetUserID(), 10, 64)

    // 构造新消息数据(包含发送者名称,避免接收方再查一次)
    newMsgData := types.NewMessageData{
        MessageID:  messageID,
        GroupID:    sendData.GroupID,
        SenderID:   senderID,
        SenderName: l.getUserName(senderID), // 带本地缓存的 RPC 调用
        MsgType:    sendData.MsgType,
        Content:    sendData.Content,
        ImageURL:   sendData.ImageURL,
        CreatedAt:  now,
    }

    // ① 立即返回 ACK(~1ms,用户体验第一)
    ackData, _ := json.Marshal(types.AckData{MessageID: messageID, Success: true})
    client.SendMessage(&types.WSMessage{
        Type:      types.TypeAck,
        MessageID: msg.MessageID,
        Timestamp: now,
        Data:      ackData,
    })

    // ② 发布到消息中间件(触发群聊广播,跨实例也能推)
    payload, _ := json.Marshal(newMsgData)
    if err := l.messagingClient.Publish(l.ctx, "chat.message.new", payload); err != nil {
        logx.Errorf("发布消息到中间件失败: %v", err)
        // 广播失败不影响落库,继续执行
    }

    // ③ 异步落库(推入队列,Worker 处理)
    task := &queue.SaveMessageTask{
        MessageID: messageID,
        GroupID:   sendData.GroupID,
        SenderID:  senderID,
        MsgType:   sendData.MsgType,
        Content:   sendData.Content,
        ImageURL:  sendData.ImageURL,
    }
    if err := l.svcCtx.SaveQueue.Push(task); err != nil {
        logx.Errorf("推送到保存队列失败(队列满): %v", err)
        // TODO: 告警通知,人工介入
    }

    return nil
}

💡 getUserName 带缓存:发送者名称通过 User RPC 查询,为避免每条消息都触发 RPC,MessageLogic 内置了一个 TTL 5 分钟的本地缓存(cache.UserCache),缓存命中率接近 100%。

异步消息保存队列(Worker Pool + 死信队列)
// 📁 app/chat/ws/internal/queue/save_queue.go

// SaveQueue 消息异步保存队列
type SaveQueue struct {
    queue          chan *SaveMessageTask // 主队列(容量 10000)
    deadLetterChan chan *SaveMessageTask // 死信队列(容量 1000)
    chatRpc        chatservice.ChatService
    wg             sync.WaitGroup
    ctx            context.Context
    cancel         context.CancelFunc
}

// NewSaveQueue 创建队列并启动 Worker Pool
func NewSaveQueue(chatRpc chatservice.ChatService, workerCount int) *SaveQueue {
    ctx, cancel := context.WithCancel(context.Background())
    sq := &SaveQueue{
        queue:          make(chan *SaveMessageTask, 10000),
        deadLetterChan: make(chan *SaveMessageTask, 1000),
        chatRpc:        chatRpc,
        ctx:            ctx,
        cancel:         cancel,
    }

    // 启动 N 个 Worker goroutine(默认 10 个)
    for i := 0; i < workerCount; i++ {
        sq.wg.Add(1)
        go sq.worker(i)
    }

    // 死信队列处理器(单独 goroutine)
    sq.wg.Add(1)
    go sq.deadLetterHandler()

    return sq
}

// processTask 处理单条保存任务(含指数退避重试)
func (sq *SaveQueue) processTask(task *SaveMessageTask) {
    _, err := sq.chatRpc.SaveMessage(sq.ctx, &chat.SaveMessageReq{
        MessageId: task.MessageID,
        GroupId:   task.GroupID,
        SenderId:  task.SenderID,
        MsgType:   task.MsgType,
        Content:   task.Content,
        ImageUrl:  task.ImageURL,
    })

    if err != nil {
        logx.Errorf("保存消息失败 (重试 %d/3): id=%s, err=%v",
            task.Retry, task.MessageID, err)

        if task.Retry < 3 {
            task.Retry++
            // 指数退避:第 1 次等 2s,第 2 次等 4s,第 3 次等 8s
            backoff := time.Duration(1<<task.Retry) * time.Second
            time.Sleep(backoff)
            sq.queue <- task // 重新入队
        } else {
            // 超过最大重试次数,进入死信队列
            logx.Errorf("消息保存失败已达上限,转入死信队列: id=%s", task.MessageID)
            sq.deadLetterChan <- task
        }
        return
    }

    logx.Infof("消息保存成功: id=%s", task.MessageID)
}

// deadLetterHandler 死信队列处理:记录日志 + 告警(可扩展)
func (sq *SaveQueue) deadLetterHandler() {
    defer sq.wg.Done()
    for {
        select {
        case <-sq.ctx.Done():
            return
        case task := <-sq.deadLetterChan:
            data, _ := json.Marshal(task)
            logx.Errorf("【死信队列】消息最终保存失败,需人工处理: %s", string(data))
            // 可扩展:写入失败消息表、发钉钉告警等
        }
    }
}

// Stop 优雅停止:等待所有任务处理完毕(最多 30 秒)
func (sq *SaveQueue) Stop() {
    sq.cancel()
    done := make(chan struct{})
    go func() {
        sq.wg.Wait()
        close(done)
    }()
    select {
    case <-done:
        logx.Info("SaveQueue 所有任务处理完成")
    case <-time.After(30 * time.Second):
        logx.Error("SaveQueue 等待超时,强制关闭")
    }
}

Worker Pool 的处理模型:

成功

失败

< 3次

= 3次

Push 入队
~1μs

主队列
cap=10000

Worker 1

Worker 2

Worker N

✅ 保存到 MySQL

重试次数

指数退避后重入队

死信队列
cap=1000

告警/人工处理

消息中间件:跨实例广播的关键

Hub 不仅向本实例的 Client 广播,还通过订阅 Redis Stream 接收其他实例发布的消息:

// 📁 app/chat/ws/hub/hub.go

// subscribeMessages 订阅消息中间件,接收跨实例推送
func (h *Hub) subscribeMessages(ctx context.Context) {
    // 订阅群聊新消息
    h.messagingClient.Subscribe("chat.message.new", "ws-message-handler",
        func(msg *message.Message) error {
            var msgData types.NewMessageData
            if err := json.Unmarshal(msg.Payload, &msgData); err != nil {
                return messaging.NewNonRetryableError(err) // 格式错误,不重试
            }

            // 广播给本实例内订阅了该群聊的所有 Client
            wsMsg := &types.WSMessage{
                Type:      types.TypeNewMessage,
                MessageID: msgData.MessageID,
                Timestamp: msgData.CreatedAt,
                Data:      json.RawMessage(msg.Payload),
            }
            h.BroadcastToGroup(msgData.GroupID, wsMsg)
            return nil
        })

    // 订阅群成员加入事件(用户加群后自动订阅群聊消息)
    h.messagingClient.Subscribe(messaging.TopicGroupMemberAdded, "ws-group-member-added",
        func(msg *message.Message) error {
            var event messaging.GroupMemberChangedEvent
            json.Unmarshal(msg.Payload, &event)

            h.mu.RLock()
            client, ok := h.clients[fmt.Sprintf("%d", event.UserID)]
            h.mu.RUnlock()

            if ok {
                h.AddClientToGroup(client, event.GroupID)
                logx.Infof("用户 %d 实时订阅群聊 %s", event.UserID, event.GroupID)
            }
            return nil
        })

    // 订阅认证进度通知(学生认证审核结果实时推送)
    h.messagingClient.Subscribe(messaging.TopicVerifyProgress, "ws-verify-progress-handler",
        func(msg *message.Message) error {
            var progressEvent messaging.VerifyProgressEventData
            json.Unmarshal(msg.Payload, &progressEvent)

            progressData, _ := json.Marshal(types.VerifyProgressData{
                VerifyID: progressEvent.VerifyID,
                Status:   progressEvent.Status,
                Refresh:  progressEvent.Refresh,
            })

            wsMsg := &types.WSMessage{
                Type:      types.TypeVerifyProgress,
                MessageID: fmt.Sprintf("verify_%d", progressEvent.VerifyID),
                Timestamp: time.Now().Unix(),
                Data:      progressData,
            }

            userID := strconv.FormatInt(progressEvent.UserID, 10)
            if err := h.SendToUser(userID, wsMsg); err != nil {
                if err == ErrUserNotOnline {
                    return nil // 用户不在线,跳过,不报错
                }
                return err
            }
            return nil
        })

    // 启动消息中间件客户端
    if err := h.messagingClient.Run(ctx); err != nil {
        logx.Errorf("消息中间件客户端停止: %v", err)
    }
}

💡 跨实例广播原理:每个 WebSocket 实例都订阅同一个 Redis Stream Topic。实例 A 的用户发消息,A 把消息发布到 Stream,实例 B 订阅后收到消息,B 上在线的群成员就收到了推送。这样 WebSocket 服务可以水平扩展,不必粘性路由。

优雅关闭

服务关闭时,需要按顺序处理:停止接受新连接 → 等待队列处理完 → 关闭消息中间件:

// 📁 app/chat/ws/websocket.go(关键片段)

quit := make(chan os.Signal, 1)
signal.Notify(quit, os.Interrupt, syscall.SIGTERM)
<-quit

logx.Info("收到退出信号,开始优雅关闭...")
cancel() // 通知 Hub.Run 退出,停止处理新连接

// 1. 停止 HTTP Server,不再接受新的 WebSocket 升级请求
server.Shutdown(context.Background())

// 2. 等待 SaveQueue 中所有消息落库(重要!避免丢消息)
if svcCtx.SaveQueue != nil {
    logx.Info("等待消息保存队列处理完成...")
    svcCtx.SaveQueue.Stop() // 内部有 30s 超时保护
}

// 3. 关闭消息中间件客户端
svcCtx.MessagingClient.Close()

logx.Info("服务已完全关闭")

运行结果

✅ 成功场景:用户发消息,群聊实时收到

[WS] 新的 WebSocket 连接已建立 | remoteAddr=192.168.1.100:54321
[WS] 用户 12345 认证成功
[WS] 用户 12345 自动订阅了 3 个群聊 [group-A, group-B, group-C]
[WS] 消息处理完成(异步): message_id=a1b2-c3d4, group_id=group-A
[MQ] 发布消息到 chat.message.new: message_id=a1b2-c3d4
[MQ] 收到广播消息,推送给 group-A 的 8 名在线成员
[Queue] 消息保存成功: id=a1b2-c3d4, duration=23ms

📊 性能指标(本地压测,1000 并发连接)

连接建立(含握手+认证):平均 12ms
消息 ACK 延迟:平均 1.2ms
群聊广播(100成员在线):平均 3ms
消息落库(异步):平均 25ms(不影响用户感知)
SaveQueue 吞吐:~3000 条/秒(10 Worker)

❌ 异常场景:连接断开自动清理

[WS] WebSocket 异常关闭: websocket: close 1006 (abnormal closure)
[WS] 用户 12345 离线,清理 3 个群聊订阅
[Redis] user:status:12345 → is_online=false, last_offline_at=1740614400

🔥 压力场景:队列过载时的死信处理

[Queue] SaveQueue 队列已满,消息推送超时: message_id=xyz-789
[Queue] 消息保存失败 (重试 1/3): id=xyz-789, err=rpc: deadline exceeded
[Queue] 消息保存失败 (重试 2/3): id=xyz-789, err=rpc: deadline exceeded
[Queue] 消息保存失败 (重试 3/3): id=xyz-789, err=rpc: deadline exceeded
[Queue] 【死信队列】消息最终保存失败,需人工处理:
        {"message_id":"xyz-789","group_id":"group-A","sender_id":12345,...}

优化建议

🚀 性能优化

  • send channel 批量写出:WritePump 中 n := len(c.send) 一次性把积压的消息批量写出,显著减少 syscall.Write 调用次数
  • getUserName 本地缓存:用 sync.Map 或小型 LRU 缓存用户名,TTL 5 分钟,避免每条消息触发 RPC
  • Worker 数量调优NewSaveQueue(chatRpc, 10) 中的 10 根据实际 Chat RPC 的并发能力调整,压测确定最优值

🛡️ 可靠性优化

  • ACK 机制:客户端收到 ACK 后才认为消息发送成功,超时未收到可重发(消息 ID 用于去重)
  • 死信队列落库:当前死信队列只记日志,生产环境应写入专门的 failed_messages 表,支持人工重处理
  • SaveQueue 监控:暴露 /metrics 接口,定期上报 queue.lengthdead_letter.length,超阈值告警

🔍 可观测性优化

  • 连接数监控GET /stats 接口暴露 online_users 当前在线数
  • 用户状态查询GET /api/users/status 接口供其他服务查询用户在线状态
  • 结构化日志:所有日志带 message_idgroup_iduser_id,方便 grep 追踪单条消息全链路

常见问题

Q1: WebSocket 和 HTTP 服务能在同一个端口吗?

A: 可以。WebSocket 握手就是 HTTP 请求(带 Upgrade: websocket Header),可以和普通 HTTP 路由共存在同一个 http.ServeMux 上。CampusHub 支持两种部署模式:

  • 独立端口/ws 路由单独监听 8889 端口,与 chat-api 的 8003 端口分离
  • 共享端口:将 WebSocketHandler 注册到 chat-api 的路由上,减少端口数量

Q2: 多实例部署时,消息怎么保证不丢、不重复推送?

A: 当前架构通过 Redis Stream 广播,每个实例独立订阅。消息发到 Stream 后,所有实例都会收到并尝试广播,如果用户只连接在某一个实例上,其他实例的 BroadcastToGroup 会因为找不到本地 Client 而静默跳过(if !ok { return })。这是"尽力而为"的广播,不会重复推送给同一用户,也不会丢失。

Q3: 大群(1000人在线)广播性能怎么样?

A: BroadcastToGroup 遍历成员集合,对每个 Client 调用 sendRaw(非阻塞,投递到 channel)。1000 个成员的遍历耗时 < 1ms(纯内存操作)。真正的瓶颈是网络 IO,由 WritePump goroutine 各自处理,完全并行。如果群规模超过 10000,可以考虑分片广播(将成员按哈希分到多个 Hub 实例)。

Q4: 客户端断线重连后,期间的离线消息怎么处理?

A: 当前 WebSocket 服务只负责实时推送,不处理离线消息。CampusHub 的方案是:客户端重连认证成功后,主动通过 HTTP API 调用 GET /chat/messages?after=lastMessageId 拉取离线期间的消息(Pull 补偿)。这样 WebSocket 服务保持无状态,易于扩展。

Q5: Token 过期了怎么处理?

A: 当前实现连接建立后不会再次校验 Token 有效期。生产环境有两种处理方式:

  1. 定期 Token 刷新:客户端在 Token 快过期时重新发 auth 消息(需服务端支持重认证)
  2. 连接有效期:Hub 记录认证时间,超过 Token 有效期主动断开连接,客户端用新 Token 重连

Q6: SaveQueue 满了怎么办?

A: Push 内部使用 1 秒超时的 select,超时返回 ErrQueueFull。当前处理是记录错误日志,消息不会落库(有丢失风险)。生产环境的完整方案:

  1. 监控 GetQueueLength(),超过 80% 容量时触发告警
  2. 降级写入:队列满时同步调用 RPC 保存(牺牲延迟换可靠性)
  3. 扩容:增加 Worker 数量或升级 Chat RPC 实例

总结

✅ 本文学习成果

通过本文,我们完整实践了一个生产级 WebSocket 群聊服务的核心技术:

  • 协议升级gorilla/websocket 的 Upgrader 配置,HTTP → WebSocket 一次握手
  • 双 goroutine 模型:ReadPump 专注读,WritePump 专注写,send channel 解耦
  • Ping/Pong 保活:54 秒定时 Ping,60 秒无响应主动断开,避免"僵尸连接"
  • Hub 连接管理:channel 序列化注册/注销,sync.RWMutex 保护并发读写
  • 连接后鉴权:auth 消息 + JWT 解析,未认证只允许 ping/auth,安全可控
  • 群聊广播:groups map 维护订阅关系,广播只遍历在线成员
  • ACK 先行:立即返回 ACK,广播和落库完全异步,延迟降至 1ms 级别
  • SaveQueue:10 Worker Pool + 指数退避重试 + 死信队列,消息落库高可靠
  • Redis Stream 解耦:发布订阅实现跨实例广播,WebSocket 服务可水平扩展
  • 优雅关闭:停止新连接 → 等待队列清空 → 关闭中间件,顺序保证不丢消息

📚 延伸阅读

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐