Go实现千万级WebSocket消息推送服务,先别急着堆机器

admin 2026-08-14 05:48:41 网络安全文章 来源:ZONE.CI 全球网 0 阅读模式

文章总结: 本文探讨了Go实现千万级WebSocket消息推送服务的架构设计,指出核心难点在于连接管理、消息推送和背压处理而非单纯增加机器。关键方案包括:将链路拆分为业务服务、消息队列和网关三层;使用分片减少锁竞争;采用单写协程保证消息顺序;设置有限发送队列防止OOM;通过用户ID分区路由并处理慢客户端。建议监控连接数、队列满、写超时等指标以定位问题。 综合评分: 87 文章分类: 实战经验,解决方案


cover_image

Go 实现千万级 WebSocket 消息推送服务,先别急着堆机器

原创

go go

Go语言教程

2026年7月28日 13:29 陕西

在小说阅读器读本章

去阅读

线上连接数刚涨起来,最先报警的不是 CPU,而是推送队列。

active_conn=186420
push_queue=64
queue_full=12783
write_timeout=936
goroutines=392817

这种日志我第一眼不会去怀疑 WebSocket 协议,更不会先调大机器。queue_full 连着涨,说明消息生产速度已经超过客户端消费速度。机器再加两倍,慢客户端还是慢,内存照样被拖死。

千万级 WebSocket 服务,难点从来不是把连接升级成功,而是连接建立以后,怎么找到它、怎么写数据、写不进去怎么办。

这里说的“千万级”,通常指整个集群承载千万长连接,不是让一台 Go 服务硬扛一千万。单机能放多少连接,取决于内存、文件描述符、网卡、消息频率和平均包大小。只报一个连接数,基本没参考价值。

我一般把链路拆成三段:

业务服务 -> 消息队列 -> WebSocket 网关 -> 客户端

业务服务只负责产生消息,不直接查 WebSocket 连接。消息按照 user_id 分区,送到对应的网关节点。网关收到消息后,只查本机连接并推送。

连接表也别直接扔进一个巨大的 sync.Map 就不管了。连接频繁上线、下线、重连时,全局竞争和 GC 压力很容易冒出来。我更习惯做分片。

type Session struct {
 UserID string
 Conn   Socket
 SendQ  chan []byte
 Done   chan struct{}
}

type sessionBucket struct {
 mu    sync.RWMutex
 items map[string]*Session
}

type SessionTable struct {
 buckets [256]sessionBucket
}

func (t *SessionTable) bucket(userID string) *sessionBucket {
 var h uint32 = 2166136261
&nbsp;for&nbsp;i :=&nbsp;0; i <&nbsp;len(userID); i++ {
&nbsp; h ^=&nbsp;uint32(userID[i])
&nbsp; h *=&nbsp;16777619
&nbsp;}
&nbsp;return&nbsp;&t.buckets[h&255]
}

func&nbsp;(t *SessionTable)&nbsp;Get(userID&nbsp;string)&nbsp;(*Session,&nbsp;bool)&nbsp;{
&nbsp;b := t.bucket(userID)
&nbsp;b.mu.RLock()
&nbsp;s, ok := b.items[userID]
&nbsp;b.mu.RUnlock()
&nbsp;return&nbsp;s, ok
}

256 个分片不是什么标准答案,只是避免所有连接挤同一把锁。真正上线前,还得看锁等待、连接分布和 CPU 核数,别看到一个数字就原样复制。

每条连接只能有一个写协程,这个约束最好从代码结构上卡死。业务协程不能拿到连接后直接 Write,否则多个 goroutine 并发写同一条连接,消息顺序、帧边界和关闭流程都会变得很难看。

func&nbsp;(s *Session)&nbsp;writeLoop()&nbsp;{
&nbsp;heartbeat := time.NewTicker(25&nbsp;* time.Second)
&nbsp;defer&nbsp;heartbeat.Stop()
&nbsp;defer&nbsp;s.Conn.Close()

&nbsp;for&nbsp;{
&nbsp;&nbsp;select&nbsp;{
&nbsp;&nbsp;case&nbsp;payload := <-s.SendQ:
&nbsp; &nbsp;_ = s.Conn.SetWriteDeadline(time.Now().Add(3&nbsp;* time.Second))
&nbsp; &nbsp;if&nbsp;err := s.Conn.WriteBinary(payload); err !=&nbsp;nil&nbsp;{
&nbsp; &nbsp;&nbsp;return
&nbsp; &nbsp;}

&nbsp;&nbsp;case&nbsp;<-heartbeat.C:
&nbsp; &nbsp;_ = s.Conn.SetWriteDeadline(time.Now().Add(2&nbsp;* time.Second))
&nbsp; &nbsp;if&nbsp;err := s.Conn.WritePing(); err !=&nbsp;nil&nbsp;{
&nbsp; &nbsp;&nbsp;return
&nbsp; &nbsp;}

&nbsp;&nbsp;case&nbsp;<-s.Done:
&nbsp; &nbsp;return
&nbsp; }
&nbsp;}
}

这里最关键的不是心跳,而是 SendQ 必须有上限。

很多 WebSocket 服务写着写着就 OOM,原因并不复杂:客户端网络变慢,发送队列持续堆积。一条连接积压几百条消息看着不多,乘上几十万连接,内存很快就没了。

推送入口要明确处理背压,不能假装消息永远写得进去。

var&nbsp;ErrClientBusy = errors.New("client send queue is full")

func&nbsp;push(s *Session, body []byte)&nbsp;error&nbsp;{
&nbsp;msg := bytes.Clone(body)

&nbsp;select&nbsp;{
&nbsp;case&nbsp;s.SendQ <- msg:
&nbsp;&nbsp;return&nbsp;nil
&nbsp;default:
&nbsp;&nbsp;return&nbsp;ErrClientBusy
&nbsp;}
}

队列满了以后怎么处理,要看业务。

聊天消息可能需要进入离线存储,客户端重连后补拉;行情、在线人数这类状态消息,旧数据通常可以直接覆盖;系统通知可以重试,但必须带消息 ID,避免客户端重复展示。

最差的做法是无限等待。一个慢客户端把消费协程卡住,后面的正常用户也跟着排队。

网关层还要解决用户路由。用户连接建立后,将 user_id -> gateway_id 写入 Redis 或专门的路由服务,并设置租约。业务消息进入 MQ 时按用户分区,消费节点发现用户不在本机,就不要继续跨节点乱转。

路由关系不能只靠客户端断开时删除。进程被强杀、机器掉电时,清理逻辑根本没机会执行。网关要定期续租,过期记录让存储层自己淘汰。

真到线上排查,我会先盯这几项:

ws_active_connections
ws_push_queue_full_total
ws_write_timeout_total
ws_message_lag_seconds
process_open_fds
go_goroutines
go_memstats_heap_inuse_bytes

连接数正常,message_lag 却不断增加,问题多半在 MQ 消费或网关写出;write_timeout 集中出现在部分运营商网络,要查客户端链路;goroutine 数量只涨不降,则要看读写协程有没有在断线后退出。

千万级 WebSocket 不是某个 Go 参数调大就能解决的。连接要分片,消息要分区,发送要串行,队列要有界,慢客户端必须舍得处理。

这几个地方没卡住,机器越多,事故规模通常也越大。


免责声明:

本文所载程序、技术方法仅面向合法合规的安全研究与教学场景,旨在提升网络安全防护能力,具有明确的技术研究属性。

任何单位或个人未经授权,将本文内容用于攻击、破坏等非法用途的,由此引发的全部法律责任、民事赔偿及连带责任,均由行为人独立承担,本站不承担任何连带责任。

本站内容均为技术交流与知识分享目的发布,若存在版权侵权或其他异议,请通过邮件联系处理,具体联系方式可点击页面上方的联系我

本文转载自:Go语言教程 go go《Go 实现千万级 WebSocket 消息推送服务,先别急着堆机器》

评论:0   参与:  0