一、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 的核心设计理念 — 全双工、持久连接、双向推送 — 永远不会过时。掌握分布式、并发、网络编程与工程化这四大支柱,才能在技术变迁中保持从容。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部