WebSocket 让服务端可以主动向客户端推送数据,是实时协作、在线游戏、行情推送的首选。但当连接数从百万迈向千万,单机内存、连接迁移、广播扇出、消息可靠性都会成为系统性难题。这篇文章拆解大规模 WebSocket 推送架构的每一层设计。
一、千万级连接的挑战
1.1 连接不是免费的
每一条 WebSocket 长连接都占用服务端资源,连接规模决定了架构形态:
| 资源 | 每连接开销 | 千万连接总量 |
|---|---|---|
| 内核 socket | fd + 收发缓冲 | 需调整内核参数 |
| 用户态对象 | Connection 对象、写缓冲 | 内存数十 GB 量级 |
| 心跳定时器 | 定时器/TCP keepalive | 需分桶批量管理 |
| 内存缓存 | 会话/订阅元数据 | 需分片与序列化 |
一句话:千万连接首先是内存与内核资源问题,其次才是应用逻辑问题。
1.2 内核参数调优
# /etc/sysctl.conf —— 支撑海量长连接
# 单进程可打开的文件描述符上限
fs.file-max = 2097152
# 本地端口范围(主动外连场景)
net.ipv4.ip_local_port_range = 1024 65535
# 全连接队列
net.core.somaxconn = 65535
# TCP 超时与 keepalive
net.ipv4.tcp_keepalive_time = 7200
net.ipv4.tcp_fin_timeout = 15
# 允许更多 TIME_WAIT 复用
net.ipv4.tcp_tw_reuse = 1
# 单进程软/硬限制
# /etc/security/limits.conf
# * soft nofile 1048576
# * hard nofile 1048576
二、连接会话管理
2.1 连接对象的内存布局
单机连接管理的关键是紧凑的内存布局 + O(1) 查找:
// Go 中的连接会话(内存友好设计)
type Session struct {
Conn *websocket.Conn
UserID string
GroupID string // 所属频道
LastBeat int64 // 最后心跳(unix 秒),用于超时清理
WriteCh chan []byte // 每连接写缓冲
// 避免为每个连接重复分配,使用 sync.Pool 复用读写缓冲
}
// O(1) 路由:userID → Session
var sessions = map[string]*Session{}
var sessionsMu sync.RWMutex
2.2 连接注册表设计
方案一:单机内存 Map
适用:单机 < 5 万连接
问题:扩容需迁移,无持久化
方案二:Redis 集中注册表
userID → { nodeID, connID }
适用:跨节点定位连接所在节点
问题:Redis 需高可用;写放大
方案三:分片本地 + 全局目录
本地节点保存完整 Session,全局目录(Redis/一致性哈希)只存路由
适用:千万级架构
// 连接建立后注册到 Redis 目录(跨节点路由用)
func registerSession(ctx context.Context, userID, nodeID, connID string) error {
key := "ws:route:" + userID
// 5 分钟 TTL + 心跳续期,节点故障后自动过期
return redisClient.SetEX(ctx, key, nodeID+"|"+connID, 5*time.Minute).Err()
}
三、网关层 WebSocket 卸载
3.1 为什么需要网关
后端业务服务通常是无状态、短请求模型,而 WebSocket 是有状态的百万长连接。把连接收敛在网关层,业务服务才能保持无状态伸缩:
客户端 ──► WebSocket 网关集群(保存连接)
│ 上行:解析消息 → 投递到消息总线
│ 下行:从消息总线订阅 → 推送给对应连接
▼
业务服务(无状态,处理业务逻辑)
▼
消息总线(Kafka/Redis Stream,广播与扇出)
3.2 网关的关键职责
| 职责 | 说明 | 示例 |
|---|---|---|
| 连接终结 | 保存 socket 与协议状态 | Nginx/Envoy/自研网关 |
| 认证握手 | 升级前校验 Token | Sec-WebSocket-Protocol 子协议 |
| 心跳代理 | 代后端维持心跳 | 网关自动回 Pong |
| 协议转换 | WS ↔ 内部消息 | 网关转 Kafka 消息 |
| 路由分发 | 把消息路由到目标连接 | userID → nodeID |
# Nginx 作为 WebSocket 卸载层:只终结 TLS 与升级
map $http_upgrade $connection_upgrade {
default upgrade;
'' close;
}
upstream ws_cluster {
server 10.0.0.11:9000;
server 10.0.0.12:9000;
server 10.0.0.13:9000;
# 长连接节点:一致性哈希保证同一用户固定同一节点
hash $arg_user_id consistent;
}
server {
listen 443 ssl http2;
server_name push.example.com;
ssl_certificate /etc/nginx/tls/fullchain.pem;
ssl_certificate_key /etc/nginx/tls/privkey.pem;
location /ws/ {
proxy_pass http://ws_cluster;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection $connection_upgrade;
proxy_set_header X-User-ID $arg_user_id;
proxy_read_timeout 3600s; # 长连接不因空闲被切断
proxy_send_timeout 3600s;
proxy_connect_timeout 5s;
}
}
一句话:网关负责「扛住连接」,业务负责「处理消息」——两者解耦后才能独立扩缩容。
四、消息广播与订阅模型
4.1 三种推送粒度
单播(P2P):一条消息只推给 1 个用户
场景:私聊、订单状态、通知
组播(Group/频道):推给一个群/频道的所有在线用户
场景:聊天室、股票行情、直播弹幕
广播(Broadcast):推给全部在线用户
场景:全站公告、全局配置下发
4.2 订阅模型的数据结构
# Redis 维护 组 → 用户 的订阅关系(set)
# 上线:SADD ws:group:room-1001 user_123
# 下线:SREM ws:group:room-1001 user_123
# 广播:SMEMBERS ws:group:room-1001 → 逐个投递
# 用户 → 组 的反向索引(用于清理)
# SADD ws:user:user_123 room-1001
4.3 广播扇出(Fan-out)架构
千万连接的广播必须在网关层扇出,而不是由业务层逐个推:
业务服务产生一条消息
│ 发布到消息总线(1 次)
▼
消息总线(Kafka topic: ws_push)
│ 每个网关节点作为消费组订阅
▼
网关节点 1 ──► 本节点持有的 50 万连接
网关节点 2 ──► 本节点持有的 50 万连接
网关节点 3 ──► 本节点持有的 50 万连接
// 网关节点消费总线消息并扇出到本地连接
func (gw *Gateway) consumeAndFanout(ctx context.Context) {
for msg := range gw.kafkaReader.Messages() {
var push PushMsg
protojson.Unmarshal(msg.Value, &push)
switch push.Scope {
case ScopeBroadcast: // 全量广播
gw.sessions.Range(func(_ string, s *Session) bool {
s.Enqueue(push.Payload) // 异步写
return true
})
case ScopeGroup: // 组播
for _, s := range gw.groupIndex[push.Group] {
s.Enqueue(push.Payload)
}
case ScopeUser: // 单播
if s := gw.sessions[push.UserID]; s != nil {
s.Enqueue(push.Payload)
}
}
}
}
五、心跳与断线重连
5.1 心跳协议设计
WebSocket 协议原生提供 Ping/Pong 控制帧。服务端应主动 Ping,客户端回 Pong;一段时间无 Pong 即判定死链:
// 服务端心跳:每 30s Ping,60s 无响应则清理
func (gw *Gateway) startHeartbeat() {
go func() {
for range time.Tick(30 * time.Second) {
now := time.Now().Unix()
gw.sessions.Range(func(id string, s *Session) bool {
if now-s.LastBeat > 60 {
s.Conn.Close() // 判定死链,触发清理
gw.removeSession(s)
return true
}
if err := s.Conn.WriteControl(websocket.PingMessage,
nil, 5*time.Second); err != nil {
s.Conn.Close()
}
return true
})
}
}()
}
一句话:心跳的本质是「用协议层可控的探活取代 TCP 不可感知的静默断开」,避免僵尸连接耗尽资源。
5.2 断线重连与幂等
客户端断线重连必须处理三个问题:抖动退避、会话恢复、消息补发。
// 客户端指数退避重连(jitter)
function connect() {
const ws = new WebSocket('wss://push.example.com/ws');
let retries = 0;
ws.onclose = () => {
const delay = Math.min(1000 * 2 ** retries, 30000) +
Math.random() * 1000; // 加抖动防惊群
retries++;
setTimeout(connect, delay);
};
ws.onopen = () => { retries = 0; };
}
重连后的消息补偿:
1. 客户端携带 lastSeq(本地收到的最后消息序号)
2. 服务端从 Redis Stream 按 seq 区间补发(离线消息)
3. 客户端按 seq 去重,保证 at-least-once / exactly-once
5.3 消息序号与去重
// 服务端为每个用户维护单调递增 seq,消息带 seq 下发
type PushMsg struct {
Seq uint64 `json:"seq"`
Payload []byte `json:"payload"`
}
// 客户端把已收到的最大 seq 回传给服务端
// 断线重连时:GET /ws/offline?after=1042
// 服务端返回 (1042, 1050] 之间的离线消息
六、分布式推送可靠性
6.1 消息可靠性层次
| 级别 | 语义 | 实现 |
|---|---|---|
| At-most-once | 可能丢失 | 直接推送,丢了就丢 |
| At-least-once | 不丢失,可能重复 | 持久化 + 客户端去重 |
| Exactly-once | 不丢不重 | 持久化 + seq 确认(成本最高) |
多数推送场景用 At-least-once + 客户端 seq 去重 即可满足。
6.2 节点故障的会话迁移
故障前:user_123 连接在 node-2
│
▼ node-2 宕机
1. Redis 路由条目 TTL 过期(5 分钟心跳未续)
2. 客户端 ws.onclose 触发重连
3. DNS/LB 把新连接调度到 node-5
4. node-5 从 Redis Stream 读取 user_123 的离线消息
5. 补发完成后恢复正常推送
6.3 背压与慢消费者
某个慢客户端(如弱网手机)会导致写缓冲堆积,进而阻塞整个推送协程。必须做每连接背压:
// 每连接有界写通道 + 非阻塞入队
const maxPending = 1024
func (s *Session) Enqueue(msg []byte) bool {
select {
case s.WriteCh <- msg:
return true
default: // 缓冲已满
// 策略:丢弃并标记该连接需重连补偿 / 主动断开
s.conn.Close()
return false
}
}
// 写协程统一串行 flush,避免并发写 socket
func (s *Session) writeLoop() {
for msg := range s.WriteCh {
s.conn.SetWriteDeadline(time.Now().Add(5 * time.Second))
if err := s.conn.WriteMessage(websocket.BinaryMessage, msg); err != nil {
return
}
}
}
七、与 SSE 的场景分界
7.1 WebSocket vs SSE 对比
| 维度 | WebSocket | SSE |
|---|---|---|
| 协议 | 独立握手(101 升级),双向 | 普通 HTTP,服务端单向推送 |
| 方向 | 全双工 | 服务端 → 客户端 |
| 二进制 | 支持 | 仅文本(可 base64 折中) |
| 自动重连 | 需自己实现 | EventSource 内置 |
| 网关兼容 | 需 Upgrade 支持 | 普通 CDN/HTTP 即可 |
| 典型场景 | 游戏、IM、协同 | 通知流、AI 流式输出 |
7.2 分界建议
用 SSE 当:
- 只需要服务端推送(通知、进度、AI token 流)
- 需要穿透简单 CDN/代理
- 想复用 HTTP 缓存与鉴权
用 WebSocket 当:
- 需要双向交互(IM、协同编辑、游戏)
- 需要二进制帧
- 推送频率极高,需要协议级复用
// SSE 一行接入
const es = new EventSource('/api/events?stream=ai');
es.onmessage = (e) => renderChunk(e.data);
es.onerror = () => es.close(); // 默认自动重连
一句话:能单向解决的问题别上 WebSocket——SSE 更简单、更好穿越代理;需要双向或高频再选 WebSocket。
八、总结
| 环节 | 关键设计 | 千万级要点 |
|---|---|---|
| 连接管理 | 紧凑 Session + O(1) 查找 | 内存布局、内核参数调优 |
| 网关卸载 | 网关终结连接、业务无状态 | 一致性哈希路由节点 |
| 广播扇出 | 消息总线 + 网关消费扇出 | Kafka/Redis 分区与消费组 |
| 心跳探活 | Ping/Pong + 分桶扫描 | 分桶批量心跳,避免全表遍历 |
| 断线重连 | 退避重连 + seq 补发 | 幂等与离线消息 |
| 背压控制 | 有界缓冲 + 慢消费者隔离 | 非阻塞入队、断连兜底 |
WebSocket 大规模架构的本质是分层与解耦:网关扛连接、总线扛消息、业务扛逻辑。连接是状态,消息是流量——把状态收敛到网关,把流量交给总线,再配合心跳、重连、背压三板斧,千万级连接才是可运维的现实,而不是可望不可及的指标。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。