断线不丢消息:WebSocket 长连接重连、心跳与消息队列实战
断线不丢消息:WebSocket 长连接重连、心跳与消息队列实战
一、断连即丢消息:WebSocket 长连接的工程化缺位
某实时协作表格上线后,用户反馈"切个 WiFi 再切回来,数据就对不上"。抓包发现,断网期间服务端推送了 4 条更新,客户端重连后只拿到最新一条。这事我见过太多团队栽进去——把 WebSocket 当成永远在线的管道,不设计重连、不设计消息序号、不设计兜底队列。
WebSocket 在协议层是长连接,但在真实网络里随时会断。WiFi 切换、运营商 NAT 超时、服务端发布重启、负载均衡健康检查误判,任何一项都会让连接掉线。浏览器原生WebSocket只暴露onclose与onerror,不会自动重连,也不会补发断连期间的消息。
断连不可怕,可怕的是断连后"假装还在"。用户继续操作,前端把消息发到一个已死的连接上,既不报错也不重试,用户以为生效,服务端却从未收到。这种"静默丢失"比直接报错更难排查。
工程化长连接必须解决四件事:断线自动重连、连接保活、消息不丢不重、并发可控。这四项决定了实时应用在生产环境能否可用。
二、指数退避、心跳保活与消息序号:长连接管理的底层机制
指数退避解决重连风暴。断线瞬间,所有客户端同时重连会把服务端打挂。退避策略让重连间隔随失败次数指数增长——首次 1 秒,第二次 2 秒,第三次 4 秒,封顶 30 秒。再加随机抖动(jitter),避免多个客户端同步重连。某 IM 产品未做退避,服务端重启后 5 万客户端同时重连,把网关压垮 40 秒。
心跳保活解决"半开连接"。TCP 层的连接可能已死,但操作系统未感知,应用层以为还在线。客户端定时发心跳包,服务端回应,超时未回即认定连接死亡,主动关闭触发重连。心跳间隔通常 15 到 45 秒,需小于 NAT 超时时间(多数运营商 60 到 120 秒)。
消息序号解决不丢不重。每条消息带单调递增的 seq 号,客户端与服务端各自维护已收发水位。重连后客户端把最后收到的 seq 上报,服务端从该 seq 之后补发。重复 seq 的消息在客户端去重队列里被丢弃。这套机制等价于 TCP 的 seq 加 ack,只是搬到应用层。
并发连接上限解决资源耗尽。浏览器对同一域名 WebSocket 连接数有限制(通常 6 个),单页应用若每个模块各开一条连接会迅速占满。应做连接复用:一个全局管理器分发消息到各模块,按 topic 路由。
综上,指数退避压制重连风暴,心跳检测半开连接,消息序号保证不丢不重,并发上限防止资源耗尽。这四件工程化事项,共同构成长连接可靠性的底层机制。
三、生产级 WebSocket 管理器:重连、心跳与消息队列
下面给出一个可复用的 WebSocket 管理器。它集成指数退避重连、心跳保活、消息序号去重、断连消息队列与并发兜底。
type WSState = 'idle' | 'connecting' | 'open' | 'closing' | 'reconnecting' | 'offline'; interface OutboundMessage { seq: number; payload: string; } /** * 可复用 WebSocket 管理器 * 集成:指数退避重连、心跳保活、消息序号去重、断连消息队列 */ export class WebSocketManager { private ws: WebSocket | null = null; private state: WSState = 'idle'; private reconnectAttempts = 0; private readonly maxReconnect = 8; private heartbeatTimer: ReturnType<typeof setInterval> | null = null; private lastPongAt = 0; private sendSeq = 0; private recvSeq = 0; // 去重窗口:收到重复 seq 直接丢弃,限制内存占用 private readonly dedupWindow = new Set<number>(); private readonly pendingQueue: OutboundMessage[] = []; private readonly listeners = new Map<string, Set<(data: unknown) => void>>(); constructor( private readonly url: string, private readonly opts: { heartbeatInterval?: number; // 心跳间隔,默认 20s heartbeatTimeout?: number; // 心跳超时,默认 45s baseBackoff?: number; // 初始退避,默认 1s maxBackoff?: number; // 最大退避,默认 30s } = {}, ) {} /** 建立连接,失败自动进入重连流程 */ connect(): void { if (this.state === 'connecting' || this.state === 'open') return; this.setState('connecting'); try { this.ws = new WebSocket(this.url); } catch (err) { // 构造异常(如 URL 非法)直接进入重连,避免状态卡死 console.warn('[WS] construct failed', err); this.scheduleReconnect(); return; } this.ws.onopen = () => this.onOpen(); this.ws.onclose = (e) => this.onClose(e); this.ws.onerror = (e) => console.warn('[WS] error', e); this.ws.onmessage = (e) => this.onMessage(e); } private onOpen(): void { this.setState('open'); this.reconnectAttempts = 0; this.lastPongAt = Date.now(); this.startHeartbeat(); // 重连成功后补发待发消息,带原 seq 以便服务端去重 this.flushPendingQueue(); // 上报最后收到的 seq,触发服务端补发断连期间消息 this.safeSend(JSON.stringify({ type: 'resume', lastSeq: this.recvSeq })); } private onClose(_e: CloseEvent): void { this.stopHeartbeat(); this.ws = null; if (this.state === 'offline' || this.state === 'closing') return; this.setState('reconnecting'); this.scheduleReconnect(); } /** 指数退避 + 随机抖动,避免重连风暴 */ private scheduleReconnect(): void { if (this.reconnectAttempts >= this.maxReconnect) { this.setState('offline'); console.warn('[WS] max reconnect reached, go offline'); return; } const base = this.opts.baseBackoff ?? 1000; const max = this.opts.maxBackoff ?? 30000; const exp = Math.min(max, base * 2 ** this.reconnectAttempts); const jitter = Math.random() * 0.3 * exp; // 30% 抖动 const delay = Math.round(exp + jitter); this.reconnectAttempts++; setTimeout(() => this.connect(), delay); } /** 心跳保活:定时 ping,超时未 pong 即认定半开,主动关闭触发重连 */ private startHeartbeat(): void { const interval = this.opts.heartbeatInterval ?? 20000; const timeout = this.opts.heartbeatTimeout ?? 45000; this.heartbeatTimer = setInterval(() => { if (Date.now() - this.lastPongAt > timeout) { // 半开连接,主动关闭让重连流程接管 this.ws?.close(4001, 'heartbeat timeout'); return; } this.safeSend(JSON.stringify({ type: 'ping', t: Date.now() })); }, interval); } private stopHeartbeat(): void { if (this.heartbeatTimer) { clearInterval(this.heartbeatTimer); this.heartbeatTimer = null; } } /** 收消息:处理 pong、ack、业务消息;按 seq 去重 */ private onMessage(e: MessageEvent): void { let data: any; try { data = JSON.parse(e.data); } catch (err) { // 非 JSON 或非法帧,丢弃但不阻断连接 console.warn('[WS] parse failed', err); return; } if (data.type === 'pong') { this.lastPongAt = Date.now(); return; } if (data.type === 'ack') { // 服务端确认收到某 seq,从待发队列移除 this.removeFromQueue(data.seq); return; } if (typeof data.seq === 'number') { // 重复 seq 直接丢弃,保证业务侧不重复处理 if (this.dedupWindow.has(data.seq)) return; this.dedupWindow.add(data.seq); // 窗口滑动,限制内存占用 if (this.dedupWindow.size > 1000) { const first = this.dedupWindow.values().next().value; if (first !== undefined) this.dedupWindow.delete(first); } this.recvSeq = Math.max(this.recvSeq, data.seq); } this.emit(data.type, data.payload); } /** * 业务侧发送消息 * 连接断开时入队,重连后补发;发送异常入队等待重连 */ send(type: string, payload: unknown): void { this.sendSeq++; const msg: OutboundMessage = { seq: this.sendSeq, payload: JSON.stringify({ type, seq: this.sendSeq, payload }), }; if (this.state !== 'open') { this.pendingQueue.push(msg); return; } this.dispatch(msg); } /** 实际写入,失败入队等待重连后补发 */ private dispatch(msg: OutboundMessage): void { if (!this.ws || this.ws.readyState !== WebSocket.OPEN) { this.pendingQueue.push(msg); return; } try { this.ws.send(msg.payload); } catch (err) { // 发送异常(如连接刚关闭)入队,等重连补发 console.warn('[WS] send failed', err); this.pendingQueue.push(msg); } } /** 重连成功后批量补发,带原 seq 以便服务端去重 */ private flushPendingQueue(): void { while (this.pendingQueue.length) { const msg = this.pendingQueue.shift()!; this.dispatch(msg); } } private removeFromQueue(seq: number): void { const idx = this.pendingQueue.findIndex(m => m.seq === seq); if (idx >= 0) this.pendingQueue.splice(idx, 1); } private safeSend(raw: string): void { if (this.ws && this.ws.readyState === WebSocket.OPEN) { try { this.ws.send(raw); } catch (err) { console.warn('[WS] safeSend', err); } } } /** 订阅消息:返回取消订阅函数,避免内存泄漏 */ on(type: string, handler: (data: unknown) => void): () => void { if (!this.listeners.has(type)) this.listeners.set(type, new Set()); this.listeners.get(type)!.add(handler); return () => this.listeners.get(type)?.delete(handler); } private emit(type: string, data: unknown): void { this.listeners.get(type)?.forEach(h => { try { h(data); } catch (err) { console.warn('[WS] listener error', err); } }); } private setState(s: WSState): void { this.state = s; this.emit('__state__', s); } /** 主动关闭:停止重连与心跳,清理队列 */ close(): void { this.setState('closing'); this.stopHeartbeat(); this.reconnectAttempts = this.maxReconnect; // 阻止后续重连 try { this.ws?.close(1000, 'normal close'); } catch (err) { console.warn('[WS] close', err); } this.ws = null; this.pendingQueue.length = 0; this.setState('offline'); } }关键点有四处。其一,重连用指数退避加随机抖动,封顶 30 秒,避免风暴。其二,心跳超时主动关闭触发重连,解决半开连接。其三,消息带 seq,客户端去重、服务端补发,保证不丢不重。其四,断连期间消息入队,重连后补发,业务侧无感知。
四、长连接的代价:资源占用、重连风暴与顺序保证陷阱
WebSocket 工程化并非无损。
第一类代价是资源占用。每条长连接在服务端占用一个文件描述符与一份内存,心跳包持续产生流量。单机几万连接是常见上限,超过需做水平扩展与连接分片。客户端常驻连接也耗电,移动端尤其敏感,需在后台时降频或断开。
第二类代价是重连风暴。即便指数退避,大规模断连仍可能压垮网关。服务端发布前应主动广播"即将重启"消息,让客户端平滑重连,配合负载均衡预热新实例。某直播平台未做主动通知,一次发布让 20 万客户端同时重连,网关丢包率飙到 12%。
第三类代价是顺序保证陷阱。消息序号保证不丢不重,但跨重连的消息顺序可能错乱——重连后服务端补发的旧消息可能晚于新消息到达。对顺序敏感的业务(如协同编辑)需在应用层做版本向量或操作转换,不能只靠 seq。
第四类代价是兼容性。部分企业代理与防火墙会掐断 WebSocket 升级握手,需降级到 HTTP 长轮询或 SSE 作为兜底。移动端 WebView 对 WebSocket 的后台保活策略各异,需做设备适配。
适用边界:实时协作、IM、推送、行情这类对延迟与消息完整性敏感的场景收益最高。低频更新、可容忍秒级延迟的场景用 HTTP 轮询更简单,不必引入长连接复杂度。
五、总结
WebSocket 长连接工程化的核心是显式管理连接生命周期、消息可靠性与重连节奏。落地建议:第一,指数退避加随机抖动做重连,封顶 30 秒避免风暴。第二,定时心跳检测半开连接,超时主动关闭触发重连。第三,消息带 seq,客户端去重、服务端补发,保证不丢不重。第四,断连期间消息入队,重连后补发,业务侧无感知。第五,顺序敏感场景在应用层做版本向量,不只靠 seq。这条路在万级并发长连接与移动端弱网场景下能跑通,回报是值得的。