ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

从零搭建实时消息推送系统:WebSocket架构、选型与排障实践

从零搭建实时消息推送系统:WebSocket架构、选型与排障实践 你肯定经历过这种场景正对着电脑敲代码QQ 或者某个网页右下角突然弹出一条消息气泡手机锁屏上微信消息一条条冒出来。这些提醒之所以能刚好在你想看到的时候出现背后靠的就是一套实时消息推送系统。我这次做的项目就是要从零搭一套类似的东西让业务端产生的事件能在几百毫秒内到达用户浏览器并弹出通知。这篇文章把整个项目的拆解、选型、编码和排障过程完整记录下来适合刚接触实时通信的后端同学以及准备给现有系统加推送能力的团队参考。1. 先把需求拆明白实时推送到底要解决什么问题1.1 一个推送请求的完整路径很多人一提推送就想到 WebSocket但实际项目最先要拆的不是技术而是链路。一条消息从产生到用户看到至少要经历四段业务系统产生事件、推送服务接收并路由、网络通道下发、客户端接收并渲染通知。我这个项目的触发场景是很典型的站内通知用户在系统里被 了、有人回复了他的评论、某个工单状态变了。这些事件本身由业务服务产生它们并不关心怎么把消息送到浏览器只负责往推送服务丢一条 JSON。推送服务要做的则是维护海量长连接把事件精准投递给对应的在线用户并且在不那么顺利的情况下保证消息不丢。这里有个容易被新手忽略的点推送服务和业务服务必须解耦。初期我是在业务代码里直接 new 一个 WebSocket 塞消息表面看能用一旦推送服务重启、升级或者连接数暴涨业务服务就被拖死了。正确的做法是业务方只往消息队列或者推送服务的 HTTP 接口投递事件推送服务内部再决定怎么下发两边互不依赖。1.2 先回答三个问题推给谁、推什么、推多少选型之前需求方给的描述往往只有一句要实时。你不把这句话拆开后面每一步都会返工。我做这个项目时先逼着自己回答了三个问题一是推给谁。是推给所有人还是按用户维度定向推送这决定了系统需不需要维护用户 ID 到连接的映射表。大多数业务都是定向的所以连接管理和路由是核心。二是推什么。是只推一个你有新消息的提示还是把完整内容也带过去如果带内容就要考虑敏感信息不能走推送通道如果不带客户端还得再拉一次接口又多了延迟和失败点。三是推多少。峰值多少条每秒、同时在线多少连接这决定了要不要上集群、要不要搞消息队列缓冲。我当时预估是单机 5 万连接、峰值 1000 TPS这个量级决定了后面很多设计的复杂度。1.3 最容易被忽略的两个隐藏需求除了功能需求有两个隐藏需求几乎每个推送系统都会遇到但很少有人一开始就考虑到。第一个是在线状态管理。推送的前提是用户在线那什么叫在线TCP 还连着不代表用户真的在看页面浏览器切后台、电脑休眠连接还在但消息已经没意义了。所以系统里必须区分连接存活和用户活跃配置里通常要加一个仅在线推送还是在线且活跃才推送的开关。第二个是客户端的多样性。同一套推送服务可能要同时服务浏览器页面、桌面客户端、甚至未来要接小程序。不同端的网络环境、断线重连策略、通知展示方式都不一样。协议设计必须统一但客户端的 SDK 要各端独立维护。我见过不少项目把所有逻辑塞在浏览器里后来要接桌面端时整段重写非常痛苦。2. 方案选型轮询、长轮询、SSE、WebSocket 到底怎么选2.1 四种主流方案横向对比聊实时推送绕不开方案选型。我给团队做过一轮对比用一张表就能看得很清楚方案实现成本实时性连接方向典型瓶颈适用场景前端定时轮询最低差取决于轮询间隔单向请求请求量大服务端压力大数据变化不频繁的后台系统长轮询中较好毫秒到秒级单向请求挂起连接连接挂起占用线程/连接资源兼容老旧环境的过渡方案SSE较低好服务端单向推送单向HTTP 协议只能服务端推客户端股票行情、日志流等单向通知WebSocket较高最好全双工双向长连接连接维护复杂需要心跳保活聊天、协作、实时互动轮询为什么会被人嫌弃其实就一个空转问题。你每 5 秒问一次服务器有消息吗绝大多数时候答案是没有但这个请求仍然消耗带宽和服务器资源。在线用户多起来之后这些无效请求会先把服务器吞掉。曾有个后台项目用 3 秒轮询300 人在线时接口 QPS 直接加了 100数据库硬生生多扛了一层无意义的压力。长轮询是当年没有 WebSocket 时的妥协方案客户端发请求后服务器不立刻返回而是挂起等有新消息再返回客户端收到响应后立刻发起下一次请求。它确实减少了空转但挂起连接对服务器的连接数压力很大而且每次传输都要重新走一遍 HTTP 头header 开销是实打实的浪费。2.2 为什么 WebSocket 成为主流选项最终我选了 WebSocket原因有三点。第一点是消息是双向的。像正在输入用户已读这类状态同步用 SSE 就做不到了因为它只能服务端推客户端。而 WebSocket 是全双工的客户端也能随时往服务端发指令一个通道解决所有问题。第二点是省报文。WebSocket 握手之后走的是帧传输一个消息帧头才几个字节相比每次 HTTP 请求动辄几百字节的 header在高频通信场景下差距非常明显。我压测时对比过同样 100 万条消息WebSocket 的服务端出口流量比轮询少了 40% 以上。第三点是生态成熟。从 Node.js 的 ws、Java 的 Netty、Go 的 gorilla/websocket到浏览器原生 WebSocket API每个环节都有稳定方案。踩坑资料也多不至于遇到诡异问题抓瞎。2.3 什么情况下别用 WebSocket虽说 WebSocket 好但它不是银弹。如果你的推送是纯单向的比如系统每 10 秒广播一次行情价格那 SSE 明显更合适——它天然支持断线重连EventSource 自带 reconnection甚至不需要额外维护心跳浏览器内核都帮你做了。再比如你就是给公司内部一个几十人用的管理后台做个订单来了的提示音那前端轮询绰绰有余完全没有引入长连接的必要。我见过一个团队为了技术先进把所有页面都改成 WebSocket结果一个后端同学要同时维护成千上万个连接运维成本瞬间拉满。还有一层现实门槛如果你的部署架构里有比较老旧的负载均衡设备、或者用了某些安全扫描设备WebSocket 的协议升级请求Upgrade 头可能被中间层拦掉。这种情况不是代码能解决的选型时就要提前确认网络链路对长连接的支持情况。3. 核心实现从零搭一个能跑起来的 WebSocket 推送服务3.1 服务端骨架Node.js ws 的最小实现技术栈我选了 Node.js一是事件驱动模型天生适合高并发长连接二是我团队里前后端都懂 JS后续要共享协议定义和 SDK 代码很方便。服务端用ws这个库它是目前 Node 生态里最成熟的 WebSocket 实现。先看核心结构。连接管理要维护两张表一张是从连接 ID 到连接对象的 Map另一张是从用户 ID 到连接 ID 的映射。为什么要两张表因为推送时需要通过用户 ID 找到连接而连接断开时需要通过连接 ID 把用户在另一张表里的记录删掉。只用一张表断开清理时你得遍历连接数过万之后就慢得没法看。const WebSocket require(ws); const wss new WebSocket.Server({ port: 8080 }); // userId - ws 连接 const userConnections new Map(); wss.on(connection, (ws, req) { // 从握手链路拿到用户身份 const userId authenticate(req); if (!userId) { ws.close(4001, unauthorized); return; } ws.userId userId; userConnections.set(userId, ws); ws.on(close, () { // 防止旧连接覆盖新连接 if (userConnections.get(userId) ws) { userConnections.delete(userId); } }); ws.on(error, (err) { console.error(connection error: ${userId}, err.message); }); }); // 业务方调这个接口推送 function pushToUser(userId, payload) { const ws userConnections.get(userId); if (ws ws.readyState WebSocket.OPEN) { ws.send(JSON.stringify(payload)); return true; } return false; // 离线交给离线流程处理 }这段代码有三处细节是实际项目中绝对不能省的。第一连接建立必须做鉴权不要依赖前端自行上报用户 ID否则任何人都能伪造身份收别人的消息。我用的是握手请求里携带的 token服务端解析后绑定到连接上。第二关闭事件里要判断 Map 里存的连接是不是当前这个防止用户多开页面时旧的连接关闭把新连接误删。第三readyState判断不能省网络异常时连接可能处于半关闭状态直接 send 会抛异常。3.2 心跳保活与断线重连小细节决定稳定性WebSocket 连接有一个经典问题TCP 连接在物理断开时服务端可能要很久才能感知到。比如用户直接合上笔记本TCP 层没有发 FIN 包服务端这个连接就一直挂在那里成为僵尸连接。解决办法是心跳。心跳的做法是服务端每 30 秒给所有连接发 ping 帧客户端收到后回 pong 帧服务端如果在 60 秒内没收到这个连接的 pong就强制关闭它并清理映射关系。ws库原生支持 ping/pong服务端代码里只需要一个定时器setInterval(() { wss.clients.forEach((ws) { if (ws.isAlive false) { ws.terminate(); return; } ws.isAlive false; ws.ping(); }); }, 30000); wss.on(connection, (ws) { ws.isAlive true; ws.on(pong, () { ws.isAlive true; }); });客户端的断线重连同样要设计好我见过最典型的问题是一断线就立刻重连形成重连风暴。服务器刚重启几百个客户端同时蜂拥重连直接把服务器又打崩一次。正确的做法是指数退避第一次断开等 1 秒重连失败等 2 秒、4 秒、8 秒最多不超过 30 秒并且每次加重连的随机抖动jitter避免多个客户端同时发起请求。function connectWithRetry() { let retry 0; function connect() { const ws new WebSocket(wss://push.example.com/ws); ws.onopen () { retry 0; }; ws.onclose () { const delay Math.min(30000, 1000 * Math.pow(2, retry)) Math.random() * 1000; retry; setTimeout(connect, delay); }; } connect(); }3.3 把推送变成屏幕右下角的系统通知服务端能推消息了接下来就是最近很多人关心的电脑右下角弹出消息这个效果。浏览器里做这件事靠的是 Notification API。注意一点这个 API 只能由用户主动交互触发授权弹窗所以要在用户点击页面按钮时先请求 permission不能等消息到了再请求否则浏览器会直接拒绝。// 页面加载后引导用户点击授权 Notification.requestPermission().then((permission) { if (permission granted) { console.log(通知权限已开启); } }); // 收到推送后弹出系统级通知 function showNotification(message) { if (Notification.permission ! granted) return; const notification new Notification(message.title, { body: message.body, icon: /logo.png, tag: message.id, // 同一 tag 的新通知会替换旧通知 requireInteraction: true, // 让通知驻留不自动消失 }); notification.onclick () { window.focus(); notification.close(); }; }这里有个容易踩的坑tag这个参数。如果你不设置 tag用户同时收到 5 条消息就会弹 5 个系统通知屏幕上直接刷屏。设置了相同 tag新通知会替换同一个旧通知。我当时的策略是聊天类的通知用会话 ID 做 tag每会话只保留最新一条任务类的通知则不用 tag因为每条都要单独提醒。还要注意页面处于后台标签页时JavaScript 的定时器可能被浏览器节流。如果依赖前端 setTimeout 做心跳后台页面容易被挂起。所以我建议心跳逻辑放在服务端主动 ping而不是依赖客户端的定时器这也是上一节服务端 ping 设计的原因之一。4. 消息可靠性发了不等于收到收到不等于看到4.1 消息 ID 与全局去重实时推送系统最容易翻车的点是消息重复。我先解释一下为什么重复无法从源头避免为了保证一条消息不丢服务端在推送超时或者客户端重连后必须重发而只要存在重发机制就必然存在客户端其实已经收到了但 ACK 确认丢失的情况于是重复消息就出现了。解决办法是给消息全局唯一 ID客户端根据 ID 去重。服务端每条消息都带一个自增或者 UUID 形式的msgId客户端维护一个最近处理过的 ID 池或者有序集合收到消息先查重再决定要不要在界面上提示。// 客户端维护最近 1000 条消息 ID const recentIds new Set(); const MAX_IDS 1000; function handlePush(message) { if (recentIds.has(message.msgId)) return; recentIds.add(message.msgId); if (recentIds.size MAX_IDS) { recentIds.delete(recentIds.values().next().value); } // 渲染与通知 renderMessage(message); }4.2 ACK 确认与重发机制纯 WebSocket 推送本质上还是尽力而为网络抖动、进程崩溃都可能让消息在传输过程中丢失。所以要想做到不丢必须加 ACK 确认。我的实现是推送服务往客户端发消息时同时把这条消息写入待确认队列Redis 里按用户存一个 pending 列表客户端收到并成功渲染后回一个{ type: ack, msgId: xxx }推送服务收到 ACK就把这条消息从待确认队列里移除。如果客户端一直没确认推送服务在连接断开或者超时后重新进入投递流程。这里有一个值得细想的设计点ACK 到底由谁发。如果客户端在渲染之前就发 ACK那 UI 弹通知那一步失败用户还是看不到如果渲染成功后才发 ACK那客户端代码就要在notification.show()成功之后才回确认。我实际采用的是展示成功算收到的策略因为推送系统的核心目标是用户看到而不是网络层送达。4.3 离线消息补发用户不在线怎么办用户离线期间产生的消息如果直接扔掉用户回来后会觉得怎么少了一堆提醒。但如果每条都存下来存储成本又很高。我的做法是分层处理临时离线掉线几分钟用 Redis 存最近 50 条消息重连后一次性补发。长时间离线则交给业务系统决定通常做法是只存你收到几条新消息的聚合提醒用户点进去再拉详情。不要把大量消息体塞进推送系统否则 Redis 内存会不堪重负。我之前就估算过如果离线 24 小时、每条消息平均 1KB、每个用户平均收 200 条一万个离线用户就是 2GB 内存这绝对不是推送系统该扛的。4.4 多实例扩展单机扛不住时的标准解法连接数超过单机上限后只能上多实例部署。但多实例带来的新问题是用户 A 连接在实例 1 上业务方随机把消息投给了实例 2实例 2 上并没有用户 A 的连接消息就丢了。标准解法是引入 Redis Pub/Sub 做跨实例广播。每个实例启动时订阅同一个 Redis channel 的推送指令业务方把消息投到 Redis 并指定目标用户所有实例都会收到这条指令再各自判断用户是否在我的连接表里是就下发不是就忽略。这样既实现了任意实例都能知道用户在哪又不需要引入额外的服务发现机制。const redis require(redis); const pub redis.createClient(); const sub redis.createClient(); // 所有实例订阅同一个 channel sub.subscribe(push:dispatch); sub.on(message, (channel, message) { const { userId, payload } JSON.parse(message); pushToUser(userId, payload); }); // 业务投递时只需要 publish function publishToUser(userId, payload) { pub.publish(push:dispatch, JSON.stringify({ userId, payload })); }这套方案的扩展能力很强但要注意 Redis 单点故障会直接导致推送停摆生产环境要给 Redis 做高可用。当实例数量超过 20 个、或者消息吞吐要求再高一些时再考虑换成消息队列如 RabbitMQ/Kafka做 fanout原理一样只是性能和可靠性更好。5. 上线前后的故障排查实录与性能参考5.1 高频问题速查表做推送系统的这几个月我整理了一份问题排查清单几乎覆盖了线上遇到过的所有怪问题先列出来现象可能原因处理方案消息发不出去服务端日志无报错用户连接已死但未被心跳清理缩短心跳周期检查僵尸连接清理逻辑客户端频繁掉线重连网络层设备超时断开空闲连接服务端心跳周期要小于设备超时时间同一用户多端登录互踢新连接覆盖旧连接确认业务是多端共存还是单端登录推送延迟高达数秒服务端事件循环被阻塞排查同步 IO 和 CPU 密集计算消息偶尔重复ACK 丢失触发了重发客户端做 msgId 去重这是最后的兜底断线后消息全部丢失未实现离线补发加 Redis 短期离线存储通知弹窗不出现页面未获授权或浏览器省电模式引导交互授权页面可见性变化时处理5.2 一次内存泄漏排查事件监听器累积上线第二周运维反馈 Node 进程内存涨到 2GB 不降。我们用heapdump抓了快照发现内存里躺着大量 WebSocket 连接对象但实际在线人数根本没那么多。逐层排查后定位到我在给process.on(message)注册监听器时放在了连接建立的回调里。每来一个连接就注册一次监听器连接关闭时监听器没移除事件监听器无限累积导致内存只涨不降。这里给所有做 Node.js 推送服务的人提个醒凡是每个连接进来时要注册的全局事件监听器一定要放在connection回调外面或者确保断开时removeListener。类似地Redis 的message事件监听器如果在连接回调内注册也会出同样的问题。5.3 Nginx 网关的坑超时和 Upgrade 头实际生产环境前面通常还有一层 Nginx 做负载均衡和 TLS 终止。WebSocket 在 Nginx 后面要做两件事一是必须设置Upgrade相关的头转发location /ws { proxy_pass http://backend_nodes; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; }如果不设置这两行Nginx 默认用 HTTP/1.0 转发请求WebSocket 握手直接失败客户端会表现成连上就断开。二是超时配置。Nginx 默认proxy_read_timeout是 60 秒意味着长连接如果 60 秒内没有数据往来Nginx 就主动掐断连接。很多团队的心跳是 30 秒一次看起来没超但如果心跳只由客户端发起、服务端没回某个瞬间超过 60 秒没有数据流连接照样断。我当时把超时调到 300 秒同时保证心跳周期远短于这个值proxy_read_timeout 300s; proxy_send_timeout 300s;5.4 压测数据与容量规划参考项目上线前我们做了一轮压测给后来人一个参考。单机 4 核 8G 的实例Node.js 版推送服务保持 5 万连接时的表现CPU 稳定在 60% 以下每秒能稳定下发 5000 条消息内存占用约 1.5GB。这个量级覆盖大部分中小业务是够的。容量规划有一个相对简单的公式单实例内存主要开销在连接对象上一条 WebSocket 连接大约消耗 5-10KB 内存包括 socket 缓冲5 万连接约 500MB消息广播的瞬时开销另算。所以 8G 内存单机撑 5 到 8 万连接比较稳妥超过这个数字就加实例不要硬扛。压测时我踩过一个误判本地压测工具和服务器在同一台机器上网络延迟为零得出的 QPS 虚高。后来把压测端放到独立机器才算数。建议任何推送系统都在真实网络环境下压测并且把心跳带来的每 30 秒一次的全量包扫描也算进 CPU 开销里。最后再分享一个个人体会实时推送系统最考验人的不是写代码那一下而是上线后那种消息到底到没到用户面前的不可见性。所以我们后来专门加了一条链路追踪日志从业务投递到服务端下发、客户端确认每一步都打一条带 messageId 的日志。排查问题时凭一个 messageId 就能把整条链路拉出来看效率提升非常明显。如果你也要做类似的系统这条建议越早落实越好别等线上出问题了再补。
返回列表