我最近因为业务上要做一套订单状态通知系统,认认真真从零搭了一遍实时消息推送系统。一开始以为只是写个接口让前端轮询就行,结果发现每分钟几千次请求的打法根本撑不住,等真正换了服务端推送才发现水比想象中深:连接管理、心跳、消息可靠投递、水平扩展、压测时各种玄学瓶颈,每个环节都能把服务搞崩。这篇就先把实时推送的原理、选型和单机实现讲清楚,再把压测踩坑和生产环境规模化的思路一并梳理出来。对于刚接触 WebSocket 推送、或者已经在用轮询方案但遇到瓶颈的读者,应该能直接拿来对照落地。
1. 先说说,为什么轮询是会呼吸的痛
1.1 一个典型场景:订单状态通知
假设你正在做一个类似外卖叫号或电商发货通知的功能,用户下单后,后端在某个时刻把订单状态从“已接单”改成“派送中”,业务要求是用户在几秒内就能在页面上看到变化。听上去很简单,于是第一版你会让前端定时去问后端“订单现在什么状态了”,也就是轮询。
轮询在并发低、状态变化慢的时候是高效的,因为实现最简单。但如果用户量上来,问题就非常明显。我那次的业务里,高峰时大概有 8000 个用户同时打开页面盯着订单列表,页面每 3 秒请求一次,光这一个接口就是 2600 多 QPS。后端每次还要查一次 Redis 和数据库来组装订单快照,整个网关带宽和数据库压力瞬间就上来了。
更要命的是延迟不均匀。你设置 3 秒轮询一次,最坏情况下要等 3 秒多才能看到新状态。用户感知还好,但你没法保证精准的“秒级通知”,万一用户手一抖刷新一下,又一个请求进来。轮询的频率一旦调到 500 毫秒,QPS 直接爆炸,服务根本扛不住。
1.2 长轮询的优化,本质上还是在“借”
后来我试过长轮询:客户端发一次请求,后端如果有状态变化就立即返回,没有变化就把连接挂住,最多挂个 20 秒再返回。这样确实显著减少了请求次数,因为 8000 个用户不再每 3 秒打一次,而是每人维持一个挂起连接,等状态变化主动返回。
但长轮询的问题在于:每个挂起来的请求都占着一个工作线程或协程,对连接本身要付出不小的内存开销,而且消息推送时是“一批请求同时被释放”,回落之后又是一窝蜂重建连接。你最终还是要维护一堆 HTTP 连接,本质上还是在用设计给“请求-响应”模式的 HTTP 去模拟真实的全双工通信,别扭不说,网关层的连接复用、超时处理都很难控制。
1.3 我需要的是服务端主动找你
最后想明白了:用户数量多了以后,真正划算的不是让客户端不停地“问”,而是让服务端在有消息时主动“推”。连接不够聪明,服务器需要能说“有事我喊你,没事你别来烦我”。这就是长连接方案的核心价值,也是我把方向转向实时消息推送系统的根本原因。
实时推送系统听起来很高大上,实际上核心就三件事:建立一个可长挂的连接,让服务端可以随时往连接上写数据,以及处理连接断开后的一系列补偿逻辑。把这三点做好了,剩下都是工程细节。
2. 别急着上 WebSocket,先看这几种方案怎么选
2.1 SSE:服务端单向推送的最优解之一
服务端推送在技术栈上不只有 WebSocket。SSE(Server-Sent Events)是一个经常被低估的方案,它使用 HTTP 协议,服务端通过text/event-stream响应持续发送数据。前端用EventSource这个 API 去接收,使用起来非常方便,还可以通过Last-Event-ID字段在断线后自动续传消息。
SSE 的优点非常明显:
- 实现简单,服务端只要把响应头设置好,不需要额外升级协议,兼容成本低。
- 自带自动重连和事件 ID,断线恢复不用自己写太多逻辑。
- 只支持服务端到客户端单向推送,协议简单不容易出错。
缺点是浏览器原生的 EventSource 只支持 GET 请求,客户端没法通过同一个连接给服务器发消息。如果业务只是“后端状态变了,前端要马上知道”,SSE 完全可以胜任,而且比 WebSocket 更轻。如果你的需求也只要单向,我建议不要用 WebSocket。
2.2 WebSocket:全双工还是得靠它
WebSocket 与 SSE 最大的区别在于全双工:客户端和服务端都可以随时往连接上写消息,不需要任何请求-响应匹配。它靠一次 HTTP Upgrade 请求把连接升级成 WebSocket 链路,之后传输的就是帧,协议头很小,开销远低于每个请求都需要大量头部信息的 HTTP。
WebSocket 帧的几种类型里,实际上日常开发只需要关心三类:
- 文本帧,服务端往客户端推业务消息都用这个。
- Ping/Pong 帧,用于心跳保活,检测死连接。
- Close 帧,主动关闭连接时使用。
它的选型逻辑很清晰:如果你需要服务器推送消息,同时客户端也要实时上报操作、主动拉取状态,或者需要做聊天室、游戏对战、协同编辑这类交互,WebSocket 是绕不开的选择。
2.3 协议选型对照
我做了个简单的对照表,方便后续做架构决策时直接参考:
| 对比项 | 普通轮询 | 长轮询 | SSE | WebSocket |
|---|---|---|---|---|
| 通信方向 | 单向请求?响应 | 单向请求?响应 | 单向服务端推 | 全双工 |
| 实时性 | 取决于轮询间隔 | 高 | 高 | 高 |
| 天然重连 | 无 | 无 | 有 | 需要自己实现 |
| 客户端支持 | 全支持 | 全支持 | 现代浏览器 | 全支持 |
| 有状态长连接开销 | 无 | 有 | 有 | 有 |
| 是否适合双向通信 | 不适合 | 不适合 | 不适合 | 非常适合 |
我个人的原则是:只通知、不交互,优先 SSE;一旦需要双向控制,直接 WebSocket。至于轮询,只在连接数很低但要兼容极老环境的场景下保留。
2.4 为什么我最终选了 WebSocket
我的业务场景里,用户不止要实时接收订单状态,还会在页面上点击“确认收货”“催单”等操作,这些操作如果通过 REST API 发,然后再走推送链路,完全可以;但既然每个用户都有一根长连接了,让这个连接同时承载上行操作,又能少一层轮询延迟,会让整体架构更顺。加上后续要支持客服在线会话,消息本来就是双向的。所以权衡下来,我最终选择以 WebSocket 为主体,REST 接口只处理登录鉴权等一次性动作。
3. 单机版服务端怎么搭,才能不翻车
3.1 先搭一个连接管理器
单机版的核心思路是维护一张“用户 ID 到 WebSocket 连接”的路由表。我用的 Node.js,WebSocket 库选的是ws,它性能足够稳,API 简单,而且自带监听 Ping/Pong 的事件,方便做心跳。
服务端骨架大概这样:
const WebSocket = require('ws'); const wss = new WebSocket.Server({ port: 8080, path: '/ws' }); const userConnections = new Map(); // userId -> Set<WebSocket> wss.on('connection', (ws) => { let userId = null; ws.on('message', (data) => { const msg = JSON.parse(data.toString()); if (msg.type === 'auth') { userId = msg.userId; if (!userConnections.has(userId)) { userConnections.set(userId, new Set()); } userConnections.get(userId).add(ws); ws.userId = userId; } else if (userId && msg.type === 'ping') { ws.send(JSON.stringify({ type: 'pong', timestamp: Date.now() })); } }); ws.on('close', () => { if (userId) { const conns = userConnections.get(userId); if (conns) { conns.delete(ws); if (conns.size === 0) userConnections.delete(userId); } } }); });为什么每个用户用一个 Set 而不是直接Map<userId, ws>?因为用户可能同时开了多个页面、手机和电脑同时在,甚至浏览器开了多个标签页。用 Set 可以天然支持同一用户多个连接,推送时逐个遍历即可。
3.2 心跳检测别偷懒,否则连接就是“僵尸”
长连接最隐蔽的问题是底层 TCP 连接断了,但服务端不一定立刻知道。如果客户端断网、突然断电、浏览器被杀,服务端只会等很多次写失败或者超时才察觉。这时候 Map 里还挂着已经没用的 WebSocket 对象,一来占内存,二来推送时你还在往死连接上写数据。
我用的是在ws库的基础上做心跳:每 30 秒向所有连接发一次 ping 帧,如果客户端在 15 秒内不回应 pong,就把这条连接打上过期标记;连续两次过期就直接调用terminate()强制关闭。
function heartbeat(ws) { if (ws.isAlive === false) { ws.terminate(); return; } ws.isAlive = false; ws.ping(); } const interval = setInterval(() => { wss.clients.forEach((ws) => { heartbeat(ws); }); }, 30000); wss.on('connection', (ws) => { ws.isAlive = true; ws.on('pong', () => { ws.isAlive = true; }); });这套机制跑起来以后,连接数从模拟期的忽高忽低变成了一个稳定曲线,内存占用也下去了。我强烈建议每个人在生产环境都加上类似心跳,它花不了几行代码,却能省掉后面大量排查死连接的时间。
3.3 消息格式与去重策略
推送的 JSON 格式要尽量固定,我用的消息结构是:
{ "msgId": "uuid-v4", "type": "order.status.changed", "userId": "10086", "content": { "orderId": "20250101001", "status": "delivering" }, "timestamp": 1704000000000 }加msgId的原因很简单:WebSocket 本身不保证消息只到达一次。我用的是 TCP,绝大多数情况不会丢,但网络抖动时客户端重连、离线消息补发,都可能导致重复消息。客户端拿到消息后按msgId去过重,比服务端简单粗暴地“只发一次”更稳妥,因为服务端永远不知道自己发出去到底有没有入浏览器脑门。
单机推送的核心代码就一句:根据userId找到连接集合,然后每个连接send(JSON.stringify(message))。
function pushToUser(userId, message) { const conns = userConnections.get(userId); if (!conns || conns.size === 0) return false; conns.forEach((ws) => { if (ws.readyState === WebSocket.OPEN) { ws.send(JSON.stringify(message)); } }); return true; }要注意readyState检查,不能只看 Map 里有连接就发,要确认连接状态确实是 OPEN,否则send会在底层触发 write 失败,完全不报错但消息丢了。这个细节我第一次写没注意,后来压测发现有些消息莫名其妙消失,排查半天就是这里。
3.4 离线消息与未读补偿
用户断网的重连期间,消息是不能丢的。所以要有一层持久化兜底:推送时如果用户不在线,就把消息写入未读存储,等用户重连上线后把未读列表拉走。业务量少时用数据库表即可;量大时建议扔给 Kafka 或 Pulsar 这类消息队列,消费端按用户消费。
这里要区分概念:实时的 WebSocket 推送只是“快车道”,可靠投递最终还得靠“慢车道”补。完整的推送流程应该是:
- 先写消息记录和未读标记(持久化)。
- 再尝试实时推送。
- 客户端收到消息后反而要主动确认,确认后再清除未读标记。
这样即使实时通道断了,用户重连后也能从历史消息拉一遍,做到最终一致。我把这套逻辑称为“先落库、后推送、再确认”,顺序不能错,否则一旦进程崩溃就会丢消息。
4. 从单机到多节点:连接分布后,问题全变了
4.1 单机再大也扛不住所有连接
单机做演示可以,上了规模就遇到天花板。一台服务器能同时维护的最大 TCP 连接数,受制于文件描述符上限、内存以及单进程事件循环的处理能力。我压测时把单机顶到 6 万个长连接,极限还能更高,但 6 万之后业务请求一来,响应延迟就开始翻倍,相当大的内存都被连接对象和内核缓冲区吃掉了。就算服务器性能再强,单点故障的脆弱性也在那儿摆着:一台机器挂了,所有用户集体断线。
生产环境必须做多节点,问题就来了:一条消息推给某个用户,这个用户的连接到底在哪个节点上?不知道。你如果只把消息发到处理业务的节点,用户连接的节点根本收不到,推送就失效。所以需要一层“让所有节点都知道连接状态”的分布式方案。
4.2 用 Redis Pub/Sub 做节点间广播,够用但不完美
最省事的方案是引入 Redis Pub/Sub:每个推送节点都订阅同一个频道,当某个节点要推送消息时,先把消息发布到 Redis 频道,所有节点都会收到这条消息,然后每个节点检查自己的本地连接表里有没有目标用户,有就发,没有就丢弃。
const redis = require('redis'); const pubClient = redis.createClient(); const subClient = pubClient.duplicate(); subClient.subscribe('push-channel'); subClient.on('message', (channel, raw) => { if (channel !== 'push-channel') return; const { userId, message } = JSON.parse(raw); pushToUser(userId, message); }); function publishToAllNodes(userId, message) { pubClient.publish('push-channel', JSON.stringify({ userId, message })); }这套设计的好处是所有节点无状态化,谁的本地连接表里有目标用户,谁就负责最终发送。但要注意一个隐含错误:如果所有节点都收到了消息,却只有目标连接所在节点能发送成功,理论上正确,但整条链路绕了一圈。消息量一旦巨大,每个节点都会白白处理很多与自己无关的消息,这就是广播风暴的雏形。
我实际这么用了一段时间后,发现高峰期 Redis 的网络流量非常高,尤其是群发、全量推送场景。于是加入了两个优化:第一,消息里带上目标节点的路由信息,比如在连接鉴权成功后,让节点把自己的节点 ID 写入 Redis 内存里“userId 在线节点”的映射;第二,推送前先查这个映射,知道目标用户在哪个节点,然后直接 Redis Pub/Sub 发布一条“定向到某节点”的消息,其他节点收到后检测 nodeId 不匹配就丢弃。虽然还是走 Pub/Sub,但每个节点的无效处理大幅减少。
4.3 在线状态与连接分布元数据
为了支撑定向路由,需要一张在线状态表。我用的方案是用 Redis Hash 存储:
- Key:
online:user:{userId}。 - Field:节点 ID。
- Field 的 Value:该节点上这个用户的连接数量。
用户连接建立时做一次HSET和计数器自增,连接断开时递减,减到 0 就删除这个 field。推送前先HGETALL拿到所有在线节点,再逐个定向发布。这样不仅支持多节点定向推送,还能顺带实现“踢下线”“全端广播”等运营能力。
这里必须处理掉线时的原子性问题:如果服务端进程崩溃,Redis 里的在线状态不会自动消失。比较稳妥的做法是让每个节点定期上报自己的心跳到 Redis(比如每隔 5 秒给在线状态加个时间戳),推送前发现某节点超过 15 秒没心跳就标记为宕机,不再向它定向推送。这块逻辑如果懒得做,最简单的替代方案是回到全量广播,用简单换正确性,也能跑,只是浪费一点流量。
4.4 为什么不用消息队列替代 Pub/Sub
有朋友问我,既然要广播,为什么不用 Kafka?事实上消息队列和 Redis Pub/Sub 是两种东西。Redis Pub/Sub 是推模式,消息发布后没有存储,消费者不在线消息就丢了;Kafka 是有持久化的订阅模型,支持消费位置记录,适合做离线消息、补偿、重放。我的经验是:实时推送的在线通道用 Redis Pub/Sub 完全够,因为它的延迟最低;而离线消息、历史记录、审计日志这类需要可靠性的东西,应该走消息队列,两者各自负责一段链路,不要混为一谈。
5. 压测做完整套骚操作,才敢说上线
5.1 压测工具怎么选
很多人会拿 ab、wrk 压一下 HTTP 接口,但压 WebSocket 长连接跟压普通接口不是一个套路。你需要工具能同时维护大量连接并周期性发送消息,常用的是websocket-bench、autobahn或者直接用 Node.js 脚本模拟几百个进程去连。我用的是自己写的一个小压测脚本,每个进程模拟固定数量的连接,连接建立后每 30 秒做一次心跳,并在特定时间触发服务端广播。
压测的指标要盯这几个:
- 最大连接数:服务器能维持多少连接不掉线。
- 消息发送延迟:从服务端 send 到客户端收到,中间耗时。
- 推送吞吐:服务端每秒能向所有连接推送多少条消息。
- 内存和 CPU 曲线:长时间运行后有没有缓慢增加的泄漏。
5.2 第一个坑:文件描述符上限
跑压测第一天就发现连接数到 1024 左右开始哗哗往下掉,新连接直接拒绝。查了下系统日志才发现进程的文件描述符到了上限 1024。Linux 下每个 TCP 连接都占一个 fd,默认的软限制太低。需要调整ulimit -n,还要在 systemd 服务里设置LimitNOFILE=1048576,否则你重启服务后配置又没生效,白折腾。这个坑每个初做长连接系统的人都会踩,别觉得是小事。
5.3 第二个坑:连接对象泄漏
压测持续 24 小时后,内存曲线没有回落,一直在涨。怀疑是 WebSocket 连接对象没被回收。查看堆快照后确认:有些连接已经不在wss.clients里,但在别的地方还有引用,比如事件监听器在连接关闭后没有被移除,或者 Redis 订阅回调里又给已经关闭的 ws 调用了 send。实际上这种问题写代码时不好发现,最好是通过周期性process.memoryUsage()记录曲线,加上对堆快照的离线分析,把长期持有引用的地方找出来,在 close 事件里统一清理监听器和相关缓存。
5.4 第三个坑:消息积压导致 OOM
广播几百条消息时没感觉,一旦一次调用循环给 5 万个连接逐个send,而客户端处理速度跟不上,消息就会在服务端发送缓冲区里堆积。Node.js 的ws.send在低位时是写内存缓冲,写不动时只回调错误。积压多了,内存直接打满,服务 OOM,进程崩溃。
应对方法是在业务代码里加背压保护:每次 send 之前检查ws.bufferedAmount,如果超过某个阈值(比如 1MB),就丢弃实时消息,让客户端通过离线消息补偿机制去拉取,而不是死等。这在推送系统里是很标准的策略:实时通道尽量轻,可靠性靠慢车道补。
5.5 最终压测数据
优化完这些问题后,我在一台 4 核 8G 的云服务器上做了最终压测,数据是这样的:
- 稳定维持 6 万个长连接,内存约 4.5GB。
- 单条消息定向推送 P99 延迟低于 20 毫秒。
- 全量广播 6 万连接同时下发一条消息,耗时约 1.8 秒。
- 断线重连场景下,客户端重连并拉取离线消息耗时 400 毫秒以内。
单机 6 万连接对绝大多数业务已经足够,而且这个数字还能往上加,只要你能接受内存和 CPU 的开销。但对生产系统来说,更重要的是有清晰的水平扩展路径,别等到 20 万连接时才开始想方案。
6. 做直播这类超高并发场景,消息量会教你做人
如果你面对的是一百万人在线观看直播这种场景,连接数和消息量不是一个量级,前面那些常规套路可能不够用。百万连接如果全走全量广播,一条弹幕要发 100 万次,就算每次都只要 1 毫秒,服务器累加起来也是极大的开销。这种时候需要按连接分组、分层推送,而不是每个连接都循环一遍。
我之前尝试过的方案是把用户按照连接所在节点和分组 ID 双重打散,消息先经过一层分发路由,路由根据连接所在节点批量写,然后把部分逻辑下沉到各节点本地。更进一步可以考虑用 UDP 组播或者 CDN 厂商的通道加速,但那已经是基础设施层面的优化了。总之,先在 1 万连接的规模把基本链路调通,再想百万并发的取舍,才是稳妥路径。
7. 最后再分享一点个人体会
实时消息推送系统做起来不难,难的是把它做得可靠和可持续扩展。我最大的收获是明白了“实时”不等于“可靠”,任何实时链路都要配一条兜底链路;轮询不是不能用,而是要在规模上意识到它的天花板。对于大多数中小团队,单机 WebSocket 加 Redis Pub/Sub 已经能覆盖 90% 的场景,千万页花太多时间在分布式推送的完美架构上,先把一条链路跑稳、把压测数据和监控做完善,才是实实在在的事。如果还有余力,于是去研究下连接状态管理、在线状态存储和推送幂等,这比盲目上 Kafka 和微服务更能提升系统的稳定性。