
聊到实时消息推送很多人第一反应就是“上WebSocket”仿佛只要把WebSocket协议一接上所有消息就能嗖嗖地飞到客户端。但等你真正把一套推送系统跑上线就会发现真正的难点根本不在握手阶段而在连接怎么管、消息怎么不丢、断线怎么重连、集群怎么扩展。这篇文章我就拿自己搭过的一套实时消息推送系统做完整拆解从方案选型一直讲到生产环境踩坑记录适合正在做技术选型或者已经在做IM、实时通知、在线协同、行情推送这类功能的同学参考。我尽量不写教科书式的东西所有结论都来自实际项目和线上故障你可以直接拿去用。1. 方案选型实时推送为什么别一上来就无脑选WebSocket先泼一盆冷水实时推送不等于WebSocket。很多场景用老掉牙的轮询反而省事选型选错了后面扩容、排障都会很痛苦。我一般是先回答三个问题再做决定数据是单向还是双向实时性要求是秒级还是毫秒级连接规模大概到什么量级1.1 短轮询和长轮询老办法为什么到今天还没退休短轮询就是前端每隔几秒发一个HTTP请求问服务端“有没有新消息”。它的优势只有一个实现极度简单任何后端语言都天然支持不需要额外库不需要专门的网关配置。但代价也明显大量请求是空转的服务端压力大实时性还受轮询间隔限制。我之前接手过一个库存提醒系统当时就是每3秒轮询一次高峰期每秒几百个请求数据库被查得冒烟改成长连接方案后QPS掉了90%以上。长轮询则是客户端发一个请求服务端挂住不返回等有数据了再响应响应完客户端立刻发下一个请求。它的实时性比短轮询好很多延迟基本在毫秒级而且兼容性好在老浏览器上也能跑。但它有老毛病连接需要反复建立和销毁服务端要保持大量挂起请求在高并发下对连接管理和超时控制非常敏感很容易把服务器线程池打满。所以我的经验是如果业务场景低频、实时性要求不苛刻、又想快速上线长短轮询完全够用。我之前做过一个报销审批通知消息一天没几条用短轮询一分钟查一次开发量几乎为零完全没必要上重家伙。1.2 WebSocket为什么是大多数实时场景的答案WebSocket是一个基于TCP的全双工通信协议一次握手建立连接后客户端和服务端可以随时互相推送消息。相比HTTP轮询它省掉了大量请求头和连接建立的重复开销在线连接状态下传输效率高得多。我记得第一次看到WebSocket的帧结构时印象最深的是它的头部极小。一个普通的数据帧如果payload在126字节以内头部只需要2字节而HTTP请求光一个请求头就要几百字节更别说每次轮询还得来一次TCP握手。这就是为什么做IM、协同编辑、实时白板这类高频双向通信时WebSocket能碾压轮询。但WebSocket也有自己的麻烦它不提供自动重连、不提供心跳保活、不提供消息确认所有这些可靠性机制都要应用层自己实现。还有就是多实例部署时连接分布在不同的机器上要向某个用户推送消息你得知道他当前连在哪个实例上这就要引入路由和状态同步机制。听完是不是觉得头大别急后面几章我专门讲怎么处理这些问题。1.3 SSE被很多人忽视的单向推送利器SSEServer-Sent Events容易被忽略但其实它在很多场景下是比WebSocket更合适的选择。它建立在HTTP之上客户端通过EventSource接口订阅服务端流式返回数据服务端可以持续推送消息给客户端。SSE最大的好处是自动重连和断线续传协议层面内置了Last-Event-ID机制断了之后重连还能从上次的位置继续接收这一点WebSocket要自己写一堆代码才能做到。另外它走普通HTTP不需要额外协议解析兼容性和网关穿透性都更好。但它的缺点是单向的客户端只能被动收不能通过同一个连接给服务端发消息而且传输格式基本限定为文本虽然可以传JSON但没法像WebSocket那样方便地传二进制数据。所以我一般这样定位如果是行情刷新、日志推送、AI对话流式输出这类单向通知业务无脑选SSE开发省一半。尤其是现在大模型流式输出带火了一阵SSE原理其实就是这么个老东西。1.4 一张表搞懂四个方案怎么选我以前做方案评审时习惯先拉一张表出来对比选型会变得非常快。这里直接把我常用的对比维度列出来维度短轮询长轮询SSEWebSocket实时性秒级~分钟级毫秒级毫秒级毫秒级通信方向单向拉取服务端推送半模拟双向服务端单向推送全双工双向服务端资源消耗高大量无效请求中高连接挂起低长连接低长连接自动重连天然支持需自研协议内置需自研消息可靠性无无有续传机制需自研ACK实现复杂度最低较低低高适用场景低频、实时性要求低老浏览器、轻量推送单向通知流、大模型输出IM、协同、双向高实时一句话总结我的选择逻辑低频看短轮询单向看SSE双向看WebSocket长轮询能不用就不用。只有当老系统实在没法升级协议、又必须提升实时性时长轮询才作为过渡方案出现。2. 架构设计拆解一个高可靠推送网关的核心链路方案定下来之后真正考验人的是把整个推送链路搭起来。我做的这套系统本质上是一个推送网关核心目标有三个连接要稳、消息要准、扩展要容易。下面我把整个架构拆开讲。2.1 一条消息从发出到上屏中间到底发生了什么假设业务方要推一条消息给用户完整链路是这样业务服务调用推送网关的HTTP接口把消息内容、接收者ID、消息类型传过来网关把这条消息写入消息队列做异步缓冲和削峰推送路由模块从在线状态表里查出这个用户当前连接在哪个网关实例上消息被转发到对应实例由该实例从连接池中取出用户对应的连接通过WebSocket推给客户端客户端收到消息后回一个ACK网关收到ACK后确认推送成功。这个链路看起来直接但每一步都有坑。比如消息队列的作用很多人不理解觉得多此一举实际上在高并发下如果没有队列做缓冲业务服务一秒钟发十万条消息网关实例可能直接被压垮反过来把推送链路打成雪崩。我做的系统里队列还承担了“消息优先发送”的功能比如IM内的消息优先于营销通知靠消息队列的优先级队列实现很顺手。2.2 连接管理Session结构、心跳与超时回收每个客户端连接到网关后服务端会建立一个会话对象我习惯叫它Session里面至少要保存这几样东西连接ID、用户ID、设备类型、客户端版本、最后活跃时间、待确认消息表。这些字段不只是为了推送更是为了排障。心跳机制是长连接最容易忽略又最容易出事的环节。TCP本身存在半开连接的问题比如客户端突然断网、手机进了隧道信号TCP连接不会立刻感知服务端还以为连接活着往这条连接上发消息就石沉大海。解决办法是应用层心跳服务端每隔一段时间主动发一个Ping帧客户端收到后回一个Pong帧如果连续几个Ping都没收到Pong就判定连接死掉了服务端主动关闭并清理Session。心跳间隔设多少有讲究。太短会增加无效流量太长会延迟感知死连接。我一般设在30秒到60秒之间具体要根据实际网络环境调。之前做过一批海外用户的服务跨网络NAT空闲超时比较短心跳就压到了25秒纯内网环境下我直接设60秒也完全没问题。2.3 消息可靠性ACK确认、超时重推与幂等去重用户最不能接受的是你告诉我消息推送成功但我根本没收到。所以推送系统的可靠性设计优先级要高于实时性。我采用的模式是“至少一次投递”网关把消息推给客户端后会把它放进该连接的待确认队列并启动一个定时器等待ACK。客户端收到消息后必须回一条ACK包网关收到ACK才把消息从待确认队列里移除。如果超过设定时间没收到ACK网关会重新推送这条消息直到成功或达到最大重试次数。这里有个很关键的问题客户端收到的消息可能是重复的。因为ACK包在网络中可能丢失服务端会重推同一条消息客户端就把同一消息显示了两遍。解决方案是客户端做幂等去重每条消息分配一个全局唯一的消息ID客户端在本地维护一个最近收到的消息ID集合收到消息后先判断ID是否已经处理过处理过就直接忽略。这里还要提一个权衡如果你追求的是“精确一次投递”那需要引入全局去重表和两阶段提交复杂度会大幅上升在绝大多数业务场景比如通知、IM中收益并不高所以我一般推荐“至少一次客户端幂等”的组合既能保证消息不丢又避免了重复带来体验问题。2.4 多实例部署时连接信息怎么共享单机推送网关撑死也就几万连接做生产系统一定要考虑多实例部署。当网关有多个节点时马上会遇到一个问题业务服务要推消息给用户但这个用户连接在哪个节点上我的方案是用Redis保存“用户在线路由表”每当用户建立连接就在Redis里写入一条映射关系key是用户IDvalue是当前网关实例ID同时设置过期时间以便崩溃时自动清理连接断开时删除这个映射。推送路由模块推送前先查一次Redis拿到用户所在实例ID再通过集群内部通道把消息转发过去。集群内部通道我用的方案是Redis的Pub/Sub。所有网关实例订阅同一个频道某个实例收到转发请求后通过频道广播出去其他实例收到广播后检查目标用户是否在本机连接池中在就推送不在就忽略。这个方案实现简单但有个扩展限制当实例数量多、广播消息量大时Pub/Sub会成为瓶颈。升级到几十个实例时建议换成基于消息队列的点对点投递或者引入gRPC做实例间直接通信路由表思路基本不变。3. 实操落地Node.js Redis搭一套带心跳和ACK的推送服务这一节我把项目里最核心的代码抽出来给一套可直接运行的思路和骨架。我选Node.js做示范因为WebSocket生态和异步模型非常适合这个场景你用Go、Java换一套语言实现核心逻辑都是同样的思路。3.1 环境准备与项目骨架我用的依赖很轻只有两个核心库ws负责WebSocket服务端ioredis负责Redis的读写和Pub/Sub。项目结构很简单push-gateway/ ├── package.json ├── src/ │ ├── server.js # 网关入口处理连接与心跳 │ ├── router.js # 在线路由表读写 │ ├── dispatcher.js # 消息分发与ACK重推 │ └── redis.js # Redis客户端封装 └── client/ └── index.html # 浏览器端联调页面npm install ws ioredis这里解释一下为什么选ws而不是socket.io。socket.io自带重连、ACK、房间机制用起来确实舒服但它抽象层次太高很多底层行为被封装得严严实实出了问题反而没法快速定位。我这种需要精细控制心跳和消息格式的场景直接用ws更顺手。你如果是快速做内部工具用socket.io也没问题只是我个人更喜欢把核心逻辑握在自己手里。3.2 服务端核心建立连接、参数鉴权与心跳保活先写一个最基础的网关入口处理连接建立和鉴权。连接建立时客户端会在URL上带一个token参数网关拿到token后先校验身份失败的直接关连接并返回错误码。// redis.js const Redis require(ioredis); // 业务缓存实例与订阅实例分开避免订阅阻塞正常读写 const businessClient new Redis({ host: 127.0.0.1, port: 6379, enableOfflineQueue: false, }); const subClient businessClient.duplicate(); subClient.on(message, handleChannelMessage); module.exports { businessClient, subClient };// server.js const { WebSocketServer } require(ws); const { businessClient } require(./redis); const { sendToUser, dispatch } require(./dispatcher); const wss new WebSocketServer({ port: 8080 }); // 连接与鉴权 wss.on(connection, async (ws, req) { const url new URL(req.url, http://localhost); const token url.searchParams.get(token); const userId url.searchParams.get(uid); // 鉴权失败直接断开 if (!token || !userId) { ws.close(4001, missing token or uid); return; } // 实际项目中这里应该调登录态校验服务 if (token ! mock-token-abc) { ws.close(4001, auth failed); return; } const connId ${userId}_${Date.now()}_${Math.random().toString(16).slice(2)}; // 保存会话信息 ws.appData { userId, connId, isAlive: true, pending: new Map() }; // 写入在线路由表TTL设置为心跳间隔的3倍 await businessClient.set(route:${userId}, gateway-1, EX, 90); // 进入业务处理 bindMessageHandler(ws); }); // 心跳检测服务端每30秒扫描一次所有连接 const HEARTBEAT_INTERVAL_MS 30 * 1000; const HEARTBEAT_TIMEOUT_MS 10 * 1000; setInterval(() { wss.clients.forEach((ws) { if (!ws.appData?.isAlive) { ws.terminate(); return; } ws.appData.isAlive false; ws.ping(); }); }, HEARTBEAT_INTERVAL_MS); // Pong响应说明连接仍然活跃 wss.on(pong, (ws) { if (ws.appData) ws.appData.isAlive true; });注意一个细节ws.terminate()和ws.close()是有区别的。close是走正常关闭流程先发关闭帧再等对端应答terminate是无条件销毁底层TCP连接对于已经判定为死掉的连接应该直接用terminate否则close可能一直卡在等待状态。3.3 消息协议设计与ACK重推机制推送消息我统一设计成JSON格式字段固定不要搞各种变体字段类型说明msgIdstring全局唯一消息ID客户端按它去重typestring消息类型message、ack、ping等payloadobject业务数据tsnumber服务端时间戳当业务方要发消息给我时我把消息塞进待确认队列启动一次性定时器超过5秒没收到ACK就重新推送。// dispatcher.js function dispatch(userId, payload) { // 模拟从路由表找到实例并推给客户端 sendToUser(userId, { msgId: createMsgId(userId), type: message, payload, ts: Date.now(), }); } function createMsgId(userId) { // 实际项目里可以用雪花算法或UUID保证全局唯一 return ${userId}-${Date.now()}-${Math.floor(Math.random() * 10000)}; }// 在server.js里补全消息处理和ACK逻辑 const PENDING_TIMEOUT_MS 5 * 1000; const MAX_RETRY_COUNT 3; function bindMessageHandler(ws) { const { userId, pending } ws.appData; ws.on(message, (raw) { let frame; try { frame JSON.parse(raw.toString()); } catch { ws.send(JSON.stringify({ type: error, payload: { code: 400, msg: bad format } })); return; } if (frame.type ack) { const msgId frame.msgId; if (pending.has(msgId)) { clearTimeout(pending.get(msgId).timer); pending.delete(msgId); } } }); ws.appData.pushMessage (msg) { let retry 0; const trySend () { if (ws.readyState ! ws.OPEN) return; ws.send(JSON.stringify(msg)); const timer setTimeout(() { if (pending.has(msg.msgId)) { retry 1; if (retry MAX_RETRY_COUNT) { trySend(); } else { pending.delete(msg.msgId); // 这里标记为“投递失败”后续走补充推送流程 } } }, PENDING_TIMEOUT_MS); pending.set(msg.msgId, { timer, retry }); }; trySend(); }; }这个重推机制看着简单但线上救了不少次。有一次我们依赖的下游推送服务返回缓慢消息在网关积累靠这套超时重推机制消息最终都补发成功没有一条真正丢失。3.4 多实例广播用Redis Pub/Sub把消息送到对的那台机器单机网关没意思我直接把多实例的发布订阅逻辑也加上。每个网关启动时都订阅同一个频道收到频道消息后检查目标用户是否在本机在就推不在就扔。// server.js底部补全订阅逻辑 const subClient require(./redis).subClient; async function handleChannelMessage(channel, message) { if (channel ! push:dispatch) return; const { userId, payload } JSON.parse(message); let targetWs null; wss.clients.forEach((ws) { if (ws.appData?.userId userId ws.readyState ws.OPEN) { targetWs ws; } }); if (targetWs) { targetWs.appData.pushMessage({ msgId: createMsgId(userId), type: message, payload, ts: Date.now(), }); } } subClient.subscribe(push:dispatch);分发路由模块的职责就是先查Redis路由表如果目标实例不是当前实例就发布到频道如果目标就在本机直接推送。这套做法最省事不需要额外维护实例列表。3.5 客户端接入代码和联调细节客户端我用原生WebSocket写联调页面主要是为了排除库封装带来的干扰。const wsUrl ws://localhost:8080/?tokenmock-token-abcuid1001; const socket new WebSocket(wsUrl); // 已处理消息缓存用于幂等去重 const processedMsgIds new Set(); socket.addEventListener(open, () { console.log(connection open); }); socket.addEventListener(message, (event) { const frame JSON.parse(event.data); if (frame.type message) { // 重复消息直接丢弃 if (processedMsgIds.has(frame.msgId)) { return; } processedMsgIds.add(frame.msgId); // 处理业务逻辑 renderMessage(frame.payload); // 回ACK socket.send(JSON.stringify({ type: ack, msgId: frame.msgId })); } if (frame.type error) { console.error(server error, frame.payload); } }); socket.addEventListener(close, (event) { console.log(connection closed, event.code, event.reason); }); socket.addEventListener(error, (err) { console.error(websocket error, err); });联调时一个小提示如果服务端没回消息先用浏览器开发者工具看WebSocket这一帧有没有发出去再在服务端打日志确认收到请求不要一上来就怀疑协议写错了90%的长连接问题都在中间链路或参数拼写上。4. 生产环境排坑断线风暴、消息乱序与内存泄漏实录代码写完只是开始真正考验系统的是上线后的各种幺蛾子。这一节我整理了三个真实踩过的大坑以及一套排查问题的方法。4.1 断线重连风暴一次服务发版把网关打垮了线上有一次我更新网关代码直接重启了服务实例。老连接瞬间全部断开客户端各自的断线重连逻辑同时触发每个客户端都在同一个时间窗口发起重连一瞬间成千上万个连接请求涌向网关结果新起来的服务又被压垮了形成恶性循环。这个问题的根源有两个一是服务端优雅下线没做好二是客户端重连策略没有加入退避和抖动。服务端这边正常的发版流程应该是先从负载均衡里摘掉该实例的流量然后通知存量连接“我要下线了”给你一个宽限期比如30秒让客户端主动断开并去连别的实例。到了宽限期还没走的连接再强杀。这样能避免瞬间断崖。客户端这边最简单的补救是指数退避随机抖动。重连间隔别用固定值第一次失败等1秒第二次2秒第三次4秒最长到60秒每次重试前再加一个0到500毫秒的随机数让重连请求散开。这个技巧能彻底化解“全网同时重连”的拥塞问题。4.2 消息重复与乱序客户端看到顺序不对怎么办消息重复的问题在第2.3节提过客户端用消息ID去重能解决。乱序问题则更隐蔽。我遇到过这样的情况用户A发一条消息给用户B两个网关实例同时把消息推出去由于网络路由原因B先收到了后面发的消息再收到前面的导致IM对话里两句话的顺序反了。排查下来发现是因为两条消息走了不同的内部转发路径到达目标实例的先后顺序不能保证。解决方案是给同一个用户的消息队列加序号机制在一个会话内生成严格递增的sequence客户端按sequence排序后渲染而不是按到达时间渲染。这样即使底层转发有延迟客户端也能恢复正确的顺序。当然代价是客户端要维护一个待排序缓冲区稍微增加了复杂度但IM场景必须这么做。4.3 连接数飙升一晚上少了三成连接是怎么回事还有一次半夜收到告警说某台网关的在线连接数比白天少了30%。我排查了半天发现连接并没有主动断开但用户消息发不过去。后来抓包才发现问题移动网络下NAT会话有超时机制长时间没有数据传输的TCP连接运营商网关会悄悄把它掐掉客户端和服务端都不知道连接看起来还“活着”实际上已经死了。这就是应用层心跳存在的意义。服务端检测到客户端超过一定时间没回Pong立即判定连接失效并清理。我还顺手做了个“客户端空闲N秒主动断开重连”的机制把这种僵尸连接从源头上控制住。排查这类问题的命令也很简单ss -s看系统整体连接状态netstat -an | grep 8080看具体端口连接数如果系统统计的TCP连接数和网关内部统计的活跃连接数对不上基本就是半开连接堆积了。4.4 常见问题速查表我挂在工位前的一份小抄这张表是我和团队挂在工位前的遇到问题先对号入座现象可能原因快速定位方法修复建议客户端收不到消息连接已死但未感知ss -s对比TCP与业务连接数缩短心跳检测周期大量连接同时断开服务端重启或网络抖动看重连日志时间分布优雅下线指数退避重连同一消息收多次ACK丢失后重推客户端打印消息ID客户端幂等去重消息顺序颠倒内部转发路径不一致打点记录sequence消息带序号客户端排序推送延迟忽高忽低Redis Pub/Sub积压监控redis-cli --stat升级点对点转发或加消费者内存持续上涨待确认队列积压抓堆dump看Map大小设置重试上限增加告警消息体过大单帧超过缓冲区看网关错误日志支持分片或限制消息大小CPU突然飙高大量JSON序列化/解析profile定位热函数优化帧格式减少大对象5. 最后说说那几条用学费换来的经验文章快结束了但我一句“综上所述”都不想写那些交过学费才攒下来的经验比任何结论都值钱。第一日志一定要打连接生命周期事件。连接建立、心跳超时、主动断开、推送失败、ACK超时每一个事件都要有结构化日志字段包含连接ID、用户ID、实例ID。没有这些日志线上出问题你都不知道从哪查起这是我吃过最大的亏。第二监控指标只盯五个就够了在线连接数、消息推送QPS、消息重推率、ACK成功率和重连频率。这五个指标能覆盖90%的推送系统健康度问题。重推率突然上升说明网络链路或客户端活跃状态出问题重连频率异常可能是服务端发版没有做优雅下线。第三上线前一定要做故障演练。找个低峰期把某个实例直接物理断网观察系统会不会因为连接大量堆积而崩溃再模拟一次推送服务整体不可用验证消息队列的积压与恢复能力。我见过太多系统平时看着挺美一断电就原形毕露。最后分享一个小技巧给每条推送消息加一个极短的生存时间TTL比如30秒后还没推成功就丢弃配合“待确认队列长度”的告警。这样既能防止消息积压撑爆内存又能及时发现问题。实时推送系统永远追求“快、准、稳”但真正的稳定不是一个协议能给你的而是靠一整套工程细节堆出来的。