一、WebSocket 协议本质与握手过程
WebSocket 是 HTML5 开始提供的一种在单个 TCP 连接上进行全双工通信的协议。与 HTTP 的请求-响应模式不同,WebSocket 允许服务端主动向客户端推送数据,真正实现了双向实时通信。
1.1 协议握手
WebSocket 连接建立始于一个 HTTP 请求。客户端发送带有特定 Header 的 HTTP 请求,服务端确认后完成协议升级:
// 客户端请求
GET /ws HTTP/1.1
Host: example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Sec-WebSocket-Version: 13
// 服务端响应
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
服务端将客户端的 Sec-WebSocket-Key 与固定 GUID "258EAFA5-E914-47DA-95CA-C5AB0DC85B11" 拼接后进行 SHA-1 哈希并 Base64 编码,生成 Sec-WebSocket-Accept 响应头完成握手验证。
1.2 数据帧结构
WebSocket 数据传输基于帧(Frame)模型,每个帧包含 FIN 标志、操作码和操作负载:
0 1 2 3
0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-------+-+-------------+-------------------------------+
|F|R|R|R| opcode|M| Payload len | Extended payload length |
|I|S|S|S| (4) |A| (7) | (16/64) |
|N|V|V|V| |S| | (if payload len==126/127) |
| |1|2|3| |K| | |
+-+-+-+-+-------+-+-------------+ - - - - - - - - - - - - - - - +
- FIN(1 bit): 是否为消息的最后一个分片
- Opcode(4 bit): 帧类型 — 0x0(继续帧), 0x1(文本), 0x2(二进制), 0x8(关闭), 0x9(Ping), 0xA(Pong)
- Mask(1 bit): 客户端到服务端的帧必须置 1(防中间代理缓存污染)
- Payload Length: 7 bit / 7+16 bit / 7+64 bit 三种长度编码
二、核心架构设计模式
2.1 发布-订阅(Pub/Sub)架构
实现百万级 WebSocket 连接的核心是解耦消息的生产者和消费者:
const Redis = require('ioredis');
const subscriber = new Redis('redis://localhost:6379');
const publisher = new Redis('redis://localhost:6379');
const userConnections = new Map();
// 每个节点订阅全局频道
subscriber.subscribe('global:broadcast');
subscriber.on('message', (channel, message) => {
const { targetUsers, data } = JSON.parse(message);
targetUsers.forEach(userId => {
const conn = userConnections.get(userId);
if (conn && conn.readyState === WebSocket.OPEN) {
conn.send(data);
}
});
});
// 业务层发布消息
function broadcastToUsers(userIds, data) {
publisher.publish('global:broadcast', JSON.stringify({
targetUsers: userIds,
data: data
}));
}
2.2 一致性哈希路由
网关层使用一致性哈希确保同一用户的请求路由到同一 WebSocket 节点:
class ConsistentHashRouter {
constructor(nodes, virtualNodes = 150) {
this.ring = new Map();
this.sortedKeys = [];
nodes.forEach(node => this.addNode(node, virtualNodes));
}
addNode(node, vNodes) {
for (let i = 0; i < vNodes; i++) {
const key = this.hash(`${node}#${i}`);
this.ring.set(key, node);
this.sortedKeys.push(key);
}
this.sortedKeys.sort((a, b) => a - b);
}
getNode(userId) {
const hash = this.hash(userId);
const idx = this.findCeiling(hash);
return this.ring.get(this.sortedKeys[idx % this.sortedKeys.length]);
}
hash(str) {
let h = 1779033703 ^ str.length;
for (let i = 0; i < str.length; i++) {
h = Math.imul(h ^ str.charCodeAt(i), 3432918353);
h = (h << 13) | (h >>> 19);
}
return h >>> 0;
}
findCeiling(hash) {
let lo = 0, hi = this.sortedKeys.length - 1;
while (lo <= hi) {
const mid = (lo + hi) >>> 1;
if (this.sortedKeys[mid] < hash) lo = mid + 1;
else hi = mid - 1;
}
return lo;
}
}
三、重连与消息可靠传输
3.1 指数退避重连策略
class ReconnectWebSocket {
constructor(url, options = {}) {
this.url = url;
this.maxRetries = options.maxRetries || 10;
this.baseDelay = options.baseDelay || 1000;
this.maxDelay = options.maxDelay || 30000;
this.jitter = options.jitter ?? 0.3;
this.retryCount = 0;
this.messageQueue = [];
this.pendingAcks = new Map();
this.connect();
}
connect() {
this.ws = new WebSocket(this.url);
this.ws.onopen = () => {
console.log('WebSocket connected, flushing queue...');
this.retryCount = 0;
this.flushQueue();
this.startHeartbeat();
};
this.ws.onclose = () => this.handleReconnect();
this.ws.onerror = (err) => console.error('WS error:', err);
this.ws.onmessage = (event) => this.handleMessage(event.data);
}
handleReconnect() {
if (this.retryCount >= this.maxRetries) {
this.onMaxRetriesReached?.();
return;
}
// 指数退避 + 随机抖动
const delay = Math.min(
this.baseDelay * Math.pow(2, this.retryCount),
this.maxDelay
);
const jitter = delay * this.jitter * (Math.random() - 0.5);
const finalDelay = Math.floor(delay + jitter);
console.log(`Reconnect attempt ${this.retryCount + 1} in ${finalDelay}ms`);
setTimeout(() => {
this.retryCount++;
this.connect();
}, finalDelay);
}
send(data) {
const message = {
id: crypto.randomUUID(),
data,
timestamp: Date.now(),
retries: 0
};
this.messageQueue.push(message);
if (this.ws?.readyState === WebSocket.OPEN) {
this.flushQueue();
}
}
flushQueue() {
while (this.messageQueue.length > 0) {
const msg = this.messageQueue[0];
this.ws.send(JSON.stringify({ id: msg.id, data: msg.data }));
this.pendingAcks.set(msg.id, { ...msg, sentAt: Date.now() });
this.messageQueue.shift();
setTimeout(() => this.checkAck(msg.id), 5000);
}
}
checkAck(msgId) {
if (this.pendingAcks.has(msgId)) {
const msg = this.pendingAcks.get(msgId);
if (msg.retries < 3) {
msg.retries++;
this.ws.send(JSON.stringify({ id: msgId, data: msg.data }));
this.pendingAcks.set(msgId, msg);
setTimeout(() => this.checkAck(msgId), 5000 * msg.retries);
} else {
this.onMessageFailed?.(msg);
this.pendingAcks.delete(msgId);
}
}
}
startHeartbeat() {
this.heartbeatTimer = setInterval(() => {
if (this.ws?.readyState === WebSocket.OPEN) {
this.ws.send(JSON.stringify({ type: 'ping', ts: Date.now() }));
}
}, 25000);
}
}
3.2 消息去重与顺序保证
class SequencedMessageHandler {
constructor() {
this.expectedSeq = 1;
this.buffer = new Map();
this.receivedIds = new Set();
this.expireTimer = null;
}
handle(packet) {
const { id, seq, data } = packet;
// 去重检查
if (this.receivedIds.has(id)) return;
this.receivedIds.add(id);
if (seq === this.expectedSeq) {
this.deliver(data);
this.expectedSeq++;
this.drainBuffer();
} else if (seq > this.expectedSeq) {
this.buffer.set(seq, data);
this.checkMissing(seq);
}
}
drainBuffer() {
while (this.buffer.has(this.expectedSeq)) {
this.deliver(this.buffer.get(this.expectedSeq));
this.buffer.delete(this.expectedSeq);
this.expectedSeq++;
}
}
checkMissing(expectedSeq) {
clearTimeout(this.expireTimer);
this.expireTimer = setTimeout(() => {
for (let s = expectedSeq; s < expectedSeq + this.buffer.size + 1; s++) {
if (!this.buffer.has(s)) {
this.requestRetransmit(s);
}
}
}, 2000);
}
deliver(data) {
// 交付给业务层处理
console.log('Delivered:', data);
}
requestRetransmit(seq) {
ws.send(JSON.stringify({ type: 'retransmit', seq }));
}
}
四、服务端高并发实现
4.1 基于 uWebSockets.js 的 Node.js 服务端
uWebSockets.js 底层使用 libuv 与自研协议栈,单节点可达百万连接,内存约 8-12 bytes/connection:
const uWS = require('uWebSockets.js');
const { v4: uuidv4 } = require('uuid');
const connectionPool = new Map();
const channelSubscribers = new Map();
const app = uWS.App();
app.ws('/ws', {
maxPayloadLength: 16 * 1024 * 1024, // 16MB
idleTimeout: 120, // 120秒无响应关闭
maxBackpressure: 64 * 1024, // 64KB 背压限制
compression: uWS.SHARED_COMPRESSOR, // 共享压缩器
open: (ws) => {
ws.userId = uuidv4();
ws.isAlive = true;
ws.subscribedChannels = new Set();
connectionPool.set(ws.userId, ws);
console.log(`[+] Connection ${ws.userId}, total: ${connectionPool.size}`);
},
message: (ws, message, isBinary) => {
const raw = Buffer.from(message).toString();
const packet = JSON.parse(raw);
switch (packet.type) {
case 'pong':
ws.isAlive = true;
break;
case 'subscribe':
handleSubscribe(ws, packet.channel);
break;
case 'broadcast':
handleBroadcast(ws, packet);
break;
}
},
close: (ws, code, message) => {
connectionPool.delete(ws.userId);
ws.subscribedChannels.forEach(ch => {
channelSubscribers.get(ch)?.delete(ws);
});
console.log(`[-] Closed ${ws.userId}, remaining: ${connectionPool.size}`);
}
});
// 心跳检测循环
setInterval(() => {
connectionPool.forEach((ws) => {
if (!ws.isAlive) return ws.terminate();
ws.isAlive = false;
ws.send(JSON.stringify({ type: 'ping', timestamp: Date.now() }));
});
}, 30000);
app.listen(9001, (token) => {
if (token) console.log('WebSocket server listening on :9001');
});
4.2 Go 语言 gorilla/websocket 实现
package main
import (
"log"
"net/http"
"sync"
"time"
"encoding/json"
"github.com/gorilla/websocket"
)
var upgrader = websocket.Upgrader{
CheckOrigin: func(r *http.Request) bool { return true },
ReadBufferSize: 4096,
WriteBufferSize: 4096,
}
type Hub struct {
connections map[*Client]bool
channels map[string]map[*Client]bool
mu sync.RWMutex
register chan *Client
unregister chan *Client
broadcast chan *Message
}
type Message struct {
Channel string `json:"channel"`
Data json.RawMessage `json:"data"`
Seq uint64 `json:"seq"`
}
type Client struct {
hub *Hub
conn *websocket.Conn
send chan []byte
channels map[string]bool
}
func NewHub() *Hub {
return &Hub{
connections: make(map[*Client]bool),
channels: make(map[string]map[*Client]bool),
register: make(chan *Client),
unregister: make(chan *Client),
broadcast: make(chan *Message, 256),
}
}
func (h *Hub) Run() {
for {
select {
case client := <-h.register:
h.mu.Lock()
h.connections[client] = true
h.mu.Unlock()
case client := <-h.unregister:
h.mu.Lock()
for ch := range client.channels {
if subs, ok := h.channels[ch]; ok {
delete(subs, client)
}
}
delete(h.connections, client)
h.mu.Unlock()
close(client.send)
case msg := <-h.broadcast:
h.mu.RLock()
clients := h.channels[msg.Channel]
h.mu.RUnlock()
for client := range clients {
select {
case client.send <- mustEncode(msg):
default:
// 慢连接保护:队列满时关闭
close(client.send)
h.mu.Lock()
delete(h.connections, client)
h.mu.Unlock()
}
}
}
}
}
func mustEncode(v interface{}) []byte {
b, _ := json.Marshal(v)
return b
}
五、分布式集群与网关方案
5.1 Nginx 负载均衡配置
使用一致性哈希实现 WebSocket 多节点的连接亲和性:
upstream websocket_backend {
# 基于 $arg_token 的一致性哈希
hash $arg_token consistent;
server ws-node-1:9001 weight=1 max_fails=2 fail_timeout=10s;
server ws-node-2:9001 weight=1 max_fails=2 fail_timeout=10s;
server ws-node-3:9001 weight=1 max_fails=2 fail_timeout=10s;
keepalive 1024;
}
server {
listen 443 ssl http2;
server_name ws.example.com;
location /ws {
proxy_pass http://websocket_backend;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
# WebSocket 超时配置
proxy_read_timeout 3600s;
proxy_send_timeout 3600s;
# 禁用缓冲实现实时性
proxy_buffering off;
proxy_buffer_size 8k;
}
}
5.2 Kong Gateway 插件方案
使用 Kong 的限流、JWT 鉴权和链路追踪插件保护 WebSocket 端点:
_format_version: "3.0"
services:
- name: websocket-service
url: http://ws-cluster:9001
routes:
- name: ws-route
paths: ["/ws"]
plugins:
- name: rate-limiting
config:
minute: 100
policy: redis
redis_host: redis
- name: jwt
config:
secret_is_base64: false
claims_to_verify: ["exp"]
- name: correlation-id
config:
header_name: X-Correlation-ID
generator: uuid
echo_downstream: true
六、性能优化与内存管理
6.1 Buffer Pool 减少 GC 压力
class BufferPool {
constructor(maxPoolSize = 10000) {
this.pool = new Array(maxPoolSize).fill(null)
.map(() => Buffer.allocUnsafe(4096));
this.free = [...Array(maxPoolSize).keys()].reverse();
this.inUse = new Set();
}
acquire() {
if (this.free.length === 0) return null;
const idx = this.free.pop();
this.inUse.add(idx);
return { idx, buf: this.pool[idx] };
}
release(idx) {
if (this.inUse.delete(idx)) {
this.free.push(idx);
}
}
}
// 全局连接控制
const MAX_CONNECTIONS_PER_NODE = 500000;
const MAX_MESSAGE_SIZE = 1024 * 1024; // 1MB
const BACKPRESSURE_THRESHOLD = 64 * 1024; // 64KB
const HEARTBEAT_INTERVAL = 25000; // 25秒
6.2 permessage-deflate 压缩
// uWebSockets.js 压缩配置
app.ws('/ws', {
compression: uWS.DEDICATED_COMPRESSOR_3KB, // 或 SHARED_COMPRESSOR
// 级别选项:
// DEDICATED_COMPRESSOR_3KB / 8KB / 16KB / 32KB / 64KB / 128KB / 256KB
maxPayloadLength: 16 * 1024 * 1024,
idleTimeout: 120,
});
// WebSocket 压缩协商
// Client: Sec-WebSocket-Extensions: permessage-deflate; client_max_window_bits=15
// Server: Sec-WebSocket-Extensions: permessage-deflate;
// server_no_context_takeover; client_no_context_takeover
// no_context_takeover 表示每条消息独立压缩,减少内存占用
七、鉴权与安全
7.1 连接级 JWT 鉴权
const jwt = require('jsonwebtoken');
app.ws('/ws', {
upgrade: (res, req, context) => {
const token = req.getQuery('token');
if (!token) {
res.writeStatus('401 Unauthorized').end();
return;
}
jwt.verify(token, SECRET_KEY, (err, decoded) => {
if (err) {
res.writeStatus('403 Forbidden').end();
return;
}
res.upgrade(
{ userId: decoded.sub, role: decoded.role },
req.getHeader('sec-websocket-key'),
req.getHeader('sec-websocket-protocol'),
req.getHeader('sec-websocket-extensions'),
context
);
});
},
open: (ws) => {
console.log(`User ${ws.userId} connected with role: ${ws.role}`);
}
});
7.2 HMAC 消息签名防篡改
const crypto = require('crypto');
function signMessage(payload, secret) {
const timestamp = Date.now();
const nonce = crypto.randomBytes(16).toString('hex');
const body = JSON.stringify({ ...payload, timestamp, nonce });
const signature = crypto
.createHmac('sha256', secret)
.update(body)
.digest('hex');
return { body, signature, timestamp, nonce };
}
function verifyMessage(received, secret, windowMs = 60000) {
const { body, signature, timestamp, nonce } = received;
// 防重放
if (Math.abs(Date.now() - timestamp) > windowMs) return false;
// 防篡改 (constant-time 比较)
const expected = crypto
.createHmac('sha256', secret)
.update(body)
.digest('hex');
return crypto.timingSafeEqual(
Buffer.from(signature),
Buffer.from(expected)
);
}
八、可观测性与监控
8.1 Prometheus 指标
const client = require('prom-client');
const wsMetrics = {
totalConnections: new client.Gauge({
name: 'ws_connections_total',
help: 'Total active WebSocket connections'
}),
messagesReceived: new client.Counter({
name: 'ws_messages_received_total',
help: 'Total messages received'
}),
messagesSent: new client.Counter({
name: 'ws_messages_sent_total',
help: 'Total messages sent'
}),
messageLatency: new client.Histogram({
name: 'ws_message_latency_seconds',
help: 'Message processing latency',
buckets: [0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1]
}),
reconnectRate: new client.Counter({
name: 'ws_reconnects_total',
help: 'Client reconnection count'
}),
errorRate: new client.Counter({
name: 'ws_errors_total',
help: 'WebSocket error count',
labelNames: ['type']
})
};
// Grafana 面板配置:
// - 活跃连接数 (Stat) - 消息 QPS (Time Series)
// - 延迟 P99/P95/P50 (Histogram) - 重连率 (Stat)
// - 错误类型分布 (Pie) - 节点资源 (CPU/Mem/FD)
// - 广播分发速率 (Gauge) - 慢连接数量 (Gauge)
8.2 OpenTelemetry 链路追踪
const { trace, context } = require('@opentelemetry/api');
function sendTraced(ws, data) {
const span = trace.getActiveSpan();
const traceContext = {};
if (span) {
const ctx = span.spanContext();
traceContext.traceId = ctx.traceId;
traceContext.spanId = ctx.spanId;
traceContext.traceFlags = ctx.traceFlags;
}
ws.send(JSON.stringify({ ...data, _trace: traceContext }));
}
function handleMessage(rawData) {
const parsed = JSON.parse(rawData);
const traceCtx = parsed._trace;
const span = tracer.startSpan('ws.message.process', {
attributes: {
'messaging.system': 'websocket',
'messaging.destination': parsed.channel,
'messaging.message_id': parsed.id
},
links: traceCtx ? [{ context: traceCtx }] : undefined
});
const ctx = trace.setSpan(context.active(), span);
return context.with(ctx, () => processBusinessLogic(parsed));
}
九、生产环境工程实践总结
| 维度 | 推荐方案 | 适用场景 |
|---|---|---|
| 超大规模连接 | uWebSockets + K8s HPA + 一致性哈希 | 游戏/IM/实时推送 (10万+) |
| 协议兼容性 | Socket.IO + Engine.IO fallback | 需降级支持老浏览器 |
| 消息可靠性 | Seq序号 + ACK确认 + 离线队列 | IM/审批通知/交易指令 |
| 实时协作 | 直传模式 + CRDT + 心跳保活 | 在线文档/协同编辑 |
| 广播分发 | Redis Pub/Sub + 分片广播 | 群发通知/直播弹幕 |
| 安全鉴权 | JWT Upgrade鉴权 + HMAC消息签名 | 金融交易/敏感数据推送 |
| 网关选型 | Nginx(简单) / Kong(功能丰富) / Envoy(云原生) | 按团队技术栈选择 |
| 消息压缩 | permessage-deflate (no_context_takeover) | 文本类消息占比高时 |
优雅关闭(Graceful Shutdown)
// K8s 滚动更新时的不掉线方案
const server = app.listen(9001, token => {});
process.on('SIGTERM', () => {
console.log('SIGTERM received, starting graceful shutdown...');
// 1. 停止接受新连接
server.close();
// 2. 通知所有客户端即将断连
connectionPool.forEach((ws, userId) => {
ws.send(JSON.stringify({ type: 'server_shutdown', reconnectIn: 30 }));
ws.close(1001, 'Server shutting down');
});
// 3. 等待现有连接处理完毕(最多30秒)
setTimeout(() => {
console.log('Force exit after grace period.');
process.exit(0);
}, 30000);
// 4. 客户端收到后延迟 30 秒再重连(等待新 Pod 启动)
});
十、WebSocket 未来演进
随着 HTTP/3 和 WebTransport 的成熟,下一代实时通信正从 TCP 向 QUIC 迁移:
- WebTransport: 基于 QUIC + HTTP/3,原生支持多流复用,0-RTT 连接建立,无队头阻塞
- WebCodecs: 浏览器原生编解码器,与 WebSocket/WebTransport 配合实现低延迟音视频
- WebTransport Datagram: 类 UDP 的不可靠传输,适用于游戏和实时音视频
但 WebSocket 的核心设计理念 — 全双工、持久连接、双向推送 — 永远不会过时。掌握分布式、并发、网络编程与工程化这四大支柱,才能在技术变迁中保持从容。

发表评论 取消回复