WebSocket 实时通信实现:群聊消息推送全链路解析
WebSocket 实时通信实现:群聊消息推送全链路解析
摘要: HTTP 是无状态的,但聊天室不能靠轮询撑着。本文以 CampusHub 校园活动平台的群聊服务为例,从 WebSocket 升级握手、Hub + Client 双层连接管理、JWT 鉴权、群聊广播,到异步消息保存队列(Worker Pool + 死信队列)、Redis Stream 消息中间件解耦、在线状态管理,完整拆解一个生产级 WebSocket 服务的实现路径。所有代码均来自真实项目。
标签: WebSocket, go-zero, 微服务, 实时通信, 群聊, Go语言, 实战
分类: 后端开发
引言
在上一篇文章中,我们解决了微服务之间"怎么说话"的问题。但校园活动平台还有另一个核心需求:群聊。
用户加入活动群聊后,能实时收到其他成员的消息——这是 HTTP 请求-响应模式天然无法实现的。让我们看一下常见的选择:
| 方案 | 原理 | 延迟 | 服务器压力 | 适用场景 |
|---|---|---|---|---|
| 短轮询 | 客户端每 N 秒请求一次 | 高(N 秒延迟) | 大量无效请求 | 不推荐 |
| 长轮询 | 请求挂起直到有数据 | 中 | 连接占用多 | 简单推送 |
| SSE | 服务端单向推流 | 低 | 较低 | 仅需服务端推 |
| WebSocket | 全双工长连接 | 极低 | 低(连接复用) | 双向实时通信 |
WebSocket 一次握手,之后全双工通信,是群聊场景的不二之选。但在微服务架构中,WebSocket 有三个额外挑战:
- 鉴权:WebSocket 握手是 HTTP 请求,但不能用普通的 JWT 中间件——连接建立后才能鉴权
- 多实例广播:消息发到实例 A,接收方的连接在实例 B,怎么跨实例推送?
- 消息可靠性:连接瞬间断了,消息还要不要保存到数据库?
CampusHub 的解决思路:同进程 Hub 管理连接 + Redis Stream 跨实例广播 + 异步队列保证消息落库。
通过本文,你将学会
- ✅ 用
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 失败,说明连接已断,退出
}
}
}
}
心跳保活的时序:
📌 示例 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)
}
}
💡 并发安全的关键:
register和unregister是 channel,修改clientsmap 的操作全部在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))
}
完整的连接建立时序:
群聊消息广播
// 📁 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,写库和广播完全异步。
发送消息逻辑(异步版本)
// 📁 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 的处理模型:
消息中间件:跨实例广播的关键
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.length和dead_letter.length,超阈值告警
🔍 可观测性优化
- 连接数监控:
GET /stats接口暴露online_users当前在线数 - 用户状态查询:
GET /api/users/status接口供其他服务查询用户在线状态 - 结构化日志:所有日志带
message_id、group_id、user_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 有效期。生产环境有两种处理方式:
- 定期 Token 刷新:客户端在 Token 快过期时重新发
auth消息(需服务端支持重认证) - 连接有效期:Hub 记录认证时间,超过 Token 有效期主动断开连接,客户端用新 Token 重连
Q6: SaveQueue 满了怎么办?
A: Push 内部使用 1 秒超时的 select,超时返回 ErrQueueFull。当前处理是记录错误日志,消息不会落库(有丢失风险)。生产环境的完整方案:
- 监控
GetQueueLength(),超过 80% 容量时触发告警 - 降级写入:队列满时同步调用 RPC 保存(牺牲延迟换可靠性)
- 扩容:增加 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 服务可水平扩展
- ✅ 优雅关闭:停止新连接 → 等待队列清空 → 关闭中间件,顺序保证不丢消息
📚 延伸阅读
更多推荐



所有评论(0)