断线不丢消息:WebSocket 长连接重连、心跳与消息队列实战

📅 2026/7/23 9:22:55
断线不丢消息: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: ReturnTypetypeof setInterval | null null; private lastPongAt 0; private sendSeq 0; private recvSeq 0; // 去重窗口收到重复 seq 直接丢弃限制内存占用 private readonly dedupWindow new Setnumber(); private readonly pendingQueue: OutboundMessage[] []; private readonly listeners new Mapstring, 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。这条路在万级并发长连接与移动端弱网场景下能跑通回报是值得的。