SLG 游戏的聊天系统是玩家社交的核心——世界频道刷屏、联盟频道指挥、私聊交易,每秒可能有上百条消息。slg-go 的聊天服务用 WebSocket 做实时推送,gRPC 做跨服务调用,Hub 模式管理频道订阅。
双协议设计
graph LR
subgraph "客户端"
A["游戏客户端"]
end
subgraph "Chat 服务"
B["WebSocket Handler"] --> C["Hub"]
D["gRPC Handler"] --> C
C --> E["Channel"]
C --> F["Filter"]
C --> G["History"]
end
A -->|"WebSocket"| B
A -->|"gRPC"| D
协议
用途
特点
WebSocket
客户端实时收发消息
长连接、低延迟、支持推送
gRPC
Logic 服务调用聊天 API
短连接、可靠、用于系统消息
客户端通过 WebSocket 发送和接收消息,Logic 服务通过 gRPC 发送系统公告、战斗结果通知等。
Hub:频道管理中心 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 type Hub struct { channels map [string ]*Channel history *RecentCache filter *ChatFilter mu sync.RWMutex } type Channel struct { ChannelType uint8 ChannelID string members map [int64 ]chan <- *Message } type Message struct { MsgID int64 Channel string SenderID int64 SenderName string TargetID int64 Content string SentAt time.Time }
Hub 是整个聊天系统的核心——管理频道、维护订阅关系、分发消息。每个 Channel 用 map[int64]chan<- *Message 存储订阅者,消息通过 channel 异步推送给每个 WebSocket 连接。
频道类型 1 2 3 4 5 6 const ( ChannelWorld uint8 = 1 ChannelAlliance uint8 = 2 ChannelPrivate uint8 = 3 ChannelSystem uint8 = 4 )
世界频道所有玩家可见,联盟频道只限联盟成员,私聊只限两个玩家。GetChannel() 按需创建频道:
1 2 3 4 5 6 7 8 9 10 11 12 13 func (h *Hub) GetChannel(channelID string ) *Channel { h.mu.Lock() defer h.mu.Unlock() if ch, ok := h.channels[channelID]; ok { return ch } ch := &Channel{ ChannelID: channelID, members: make (map [int64 ]chan <- *Message), } h.channels[channelID] = ch return ch }
WebSocket 处理:读写分离 每个 WebSocket 连接启动两个 goroutine——读循环和写循环:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 func newWebSocketHandler (hub *chat.Hub, tokenSecret string ) http.Handler { return websocket.Handler(func (ws *websocket.Conn) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() msgCh := make (chan *chat.Message, 128 ) controlCh := make (chan wsEnvelope, 64 ) go func () { defer writeWG.Done() for { select { case <-ctx.Done(): return case msg := <-msgCh: websocket.JSON.Send(ws, ...) case env := <-controlCh: websocket.JSON.Send(ws, ...) } } }() for { var cmd wsCommand websocket.JSON.Receive(ws, &cmd) } }) }
读循环在主 goroutine 中执行(阻塞读),写循环在子 goroutine 中执行。msgCh 缓冲 128 条消息——如果客户端消费太慢,超出缓冲的消息会被丢弃并触发告警。
消息广播:异步分发 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 func (h *Hub) Broadcast(channelID string , msg *Message) { h.mu.RLock() ch, ok := h.channels[channelID] h.mu.RUnlock() if !ok { return } h.history.Add(channelID, msg) for _, msgCh := range ch.members { select { case msgCh <- msg: default : } } }
select + default 非阻塞发送——如果某个消费者的 channel 已满,直接丢弃而不是阻塞整个广播。这保证了慢消费者不会拖垮整个频道。
敏感词过滤:发送前拦截 消息发送前经过 ChatFilter(Aho-Corasick 自动机)过滤:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 func (h *Hub) Send(senderID int64 , channelID, content string ) (*Message, error ) { filtered := h.filter.Filter(content) if filtered != content { } msg := &Message{ MsgID: generateID(), Channel: channelID, SenderID: senderID, Content: filtered, SentAt: time.Now(), } h.Broadcast(channelID, msg) return msg, nil }
过滤在 Hub 层而非 WebSocket 层,保证无论是 WebSocket 还是 gRPC 发送的消息都会被过滤。
消息持久化:RecentCache 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 type RecentCache struct { mu sync.RWMutex messages map [string ][]*Message limit int } func (c *RecentCache) Add(channelID string , msg *Message) { c.mu.Lock() defer c.mu.Unlock() c.messages[channelID] = append (c.messages[channelID], msg) if len (c.messages[channelID]) > c.limit { c.messages[channelID] = c.messages[channelID][len (c.messages[channelID])-c.limit:] } }
内存缓存最近 1000 条消息,新玩家加入频道时可以获取历史消息。不需要持久化到数据库——重启后消息丢失对 SLG 聊天来说可以接受。
gRPC 接口:跨服务调用 Logic 服务通过 gRPC 发送系统消息:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 type ChatGrpcHandler struct { hub *chat.Hub } func (h *ChatGrpcHandler) SendSystemMessage(ctx context.Context, req *pb.SendSystemMessageRequest) (*pb.SendSystemMessageResponse, error ) { msg := &chat.Message{ Channel: req.ChannelId, SenderID: 0 , SenderName: "System" , Content: req.Content, SentAt: time.Now(), } h.hub.Broadcast(req.ChannelId, msg) return &pb.SendSystemMessageResponse{MsgId: msg.MsgID}, nil }
gRPC 接口复用同一个 Hub,保证消息格式和过滤逻辑一致。
慢消费者保护 当某个 WebSocket 连接消费太慢时,msgCh(缓冲 128 条)会溢出。Hub 的 Broadcast 方法用 select + default 丢弃溢出消息,同时触发告警:
1 2 3 4 5 6 7 8 9 10 func (h *Hub) shouldLogDropWarn(now time.Time) bool { h.mu.Lock() defer h.mu.Unlock() if now.Sub(h.lastDropWarn) < h.dropWarnCooldown { return false } h.lastDropWarn = now return true }
告警有冷却时间(默认 30 秒),避免日志被刷爆。运维可以根据告警频率判断是否需要扩容或优化客户端消费逻辑。
总结
设计决策
原因
WebSocket + gRPC 双协议
客户端用 WebSocket 实时推送,Logic 用 gRPC 系统调用
Hub 模式管理频道
中心化订阅关系,避免消息路由复杂度
channel 异步推送
非阻塞分发,慢消费者不拖垮整体
内存缓存历史消息
1000 条/频道,新玩家加入可获取上下文
Aho-Corasick 过滤
Hub 层统一过滤,WebSocket 和 gRPC 都经过
聊天系统不是”加个 WebSocket 就行”的事。频道管理、慢消费者保护、敏感词过滤、历史消息,每个环节都要考虑。