ARTICLE DETAIL

资讯详情

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

实时消息推送系统架构:从轮询到WebSocket的实践与优化

实时消息推送系统架构:从轮询到WebSocket的实践与优化 1. 技术选型我为什么把轮询换成了 WebSocket做实时消息推送系统之前我一直觉得推送这件事很简单——客户端定时拉一下接口不就行了直到业务方把需求拍到我桌上消息要在 1 秒内触达用户同时在线量要支持几十万级别而且消息不能丢。轮询方案直接出局反复权衡之后我选了一条更主流也更稳妥的路子以 WebSocket 作为长连接通道配合服务端主动推送。先说一个核心认知实时推送的本质是反向请求。传统 HTTP 接口是客户端主动发起请求、服务端返回结果而推送系统要求服务端在某个业务事件发生时主动把数据塞给客户端。轮询能做到这件事但代价巨大——每次请求都要携带完整 HTTP 头服务端要处理大量无效请求而且轮询间隔越短实时性越好服务端压力也越大。实测下来一个 10 万在线的业务如果轮询间隔是 3 秒网关层每秒要扛 3 万多 HTTP 请求其中绝大多数是无更新时的空转。对比之下WebSocket 只建立一次 TCP 连接后续所有消息都走这条全双工通道没有重复握手更没有无效轮询。1.1 三种实时方案的真实对比我在项目启动前做了一个对比表格把轮询、SSEServer-Sent Events和 WebSocket 放在一起看心里就非常有数了对比维度短轮询SSEWebSocket传输方向客户端拉取服务端单向推送全双工双向传输实时性取决于轮询间隔秒级有断线重试机制毫秒级连接开销每次请求都走完整 HTTP长连接但头部开销高长连接头部开销低断线重连天然支持内置自动重连需要自己实现二进制数据任意 HTTP 响应仅文本帧文本帧和二进制帧都支持服务端实现复杂度最低较低中等需要处理心跳和粘包SSE 其实是个容易被低估的方案——如果业务只需要服务端单向推送比如股票价格刷新、公告通知、日志流SSE 的自动重连优势非常省心。但我们既然叫实时消息推送系统往往还伴随着客户端回执、在线状态上报这类上行消息所以 WebSocket 的双向能力更合适。我自己的经验是别一上来就迷信 WebSocket。如果需求只是服务端往客户端推SSE 可以帮你省掉心跳和重连的底层工作。等需求里出现客户端需要向服务端发信令的时候再切 WebSocket 不迟。1.2 业务场景对选型的决定性影响选型不能只看技术指标还得看业务形态。我们的推送系统一开始有两类典型场景第一类是 IM 私信。消息由某个用户发起目标用户要立刻收到还需要送达回执。这类场景天然需要双向通道——发送端需要知道消息已到达接收端需要回 ACK。第二类是系统通知。比如订单状态变更、优惠券到账这类消息由服务端触发用户只是被动接收但量非常大可能一次活动就要推送几十万条。这两类场景混在一起决定了我的系统要同时具备长连接管理和离线消息容错两种能力。前者靠 WebSocket 网关后者靠消息存储和拉取接口配合。还有一点容易被忽略客户端运行的网络环境。移动端在 Wi-Fi 和蜂窝网络之间切换、App 被系统杀进程、弱网丢包严重……这些都会导致连接频繁断开。所以推送系统真正的难点不在于建立连接推一条消息而在于连接断了之后怎么把消息补上。这个我放在后面第 3 章重点讲。1.3 最终架构与接入层设计最终我的架构分为四层生产者业务服务 → 推送中间件消息队列 路由 → 接入网关WebSocket 长连接集群 → 客户端App / 浏览器接入层我单独抽出了一个网关服务没有让业务服务直接跟客户端通信。这样做的理由很实际业务服务不该关心客户端在线状态这是网关的职责多业务共用一套推送通道避免每个服务都维护一个长连接体系网关水平扩展时可以独立扩容不影响业务服务接入层的核心是连接注册与路由。每个客户端连接建立后网关会把连接信息连接 ID、用户 ID、网关节点地址写入一个全局路由表。推送时先查路由表找到用户落在哪个网关上再把消息转发过去。这里我用的是 Redis 存储路由信息Key 是用户 IDValue 是网关节点和连接 IDTTL 设为心跳周期的三倍防止路由数据过期残留。这一层最需要注意的坑是不要在网关内存里维护路由。单个网关只知道自己的连接如果用户连到了 A 网关推送系统却把消息发到了 B 网关B 网关查不到连接就只能丢弃。集中式路由表虽然多了一次 Redis 访问但换来的是推送路径的正确性。2. 连接生命周期管理心跳、重连与僵尸连接清理很多人以为 WebSocket 连上就能一直用实际上网络环境远比想象中脆弱Nginx 超时断开、手机锁屏后网络链路静默、代理服务器回收空闲连接……随便一个都能把连接偷偷断掉而服务端和客户端都还蒙在鼓里以为连接是好的。2.1 心跳机制的具体参数设计心跳是解决假死连接的唯一办法。所谓假死就是 TCP 连接本身还挂着但数据链路已经不可用。心跳的原理是客户端定时发一个 ping 帧服务端收到后回 pong 帧以此证明链路仍然活着。我的参数设计是这样的客户端每 30 秒发送一个应用层心跳包不是 WebSocket 协议层的 ping 帧而是自定义的 JSON 心跳消息服务端连续 3 次90 秒未收到心跳判定连接死掉主动断开并清理路由服务端自己也发心跳间隔也是 30 秒客户端如果 90 秒没收到服务端心跳主动发起重连为什么不用协议层的 ping/pong因为 WebSocket 的 ping/pong 帧在部分代理层是透明的中间可能被吞掉而应用层心跳是业务数据穿透性更好还能顺便携带客户端时间戳做延迟统计。注意心跳间隔不是越短越好。我见过有人设 5 秒心跳连接是保住了但 10 万在线每秒就有 2 万个心跳包网关压力白白增加。30 秒是平衡实时性和开销之后的经验值。2.2 僵尸连接怎么拖垮了我们的网关这个坑我印象特别深。系统刚上线那阵子网关的连接数一路走高从 5 万涨到 9 万但在线用户数其实只有 4 万多。按说连接数应该跟在线数接近多出来的近 5 万连接是哪里来的查了一圈发现大量连接是客户端已经杀掉 App 却没能正常发送关闭帧留下的残骸。TCP 连接本身不会因为进程消失而立刻消失如果客户端异常退出服务端只能靠超时判断。而当时我们没有兜底清理机制僵尸连接越积越多最终把网关的文件描述符耗光了。解决方式是双管齐下第一在服务端加了一个空转检测后台任务。每隔 5 分钟扫描所有连接凡是超过 90 秒没有任何数据收发的连接直接强关。这个逻辑独立于心跳判断专门对付那些心跳超时判不出来、但确实已经没数据的连接。第二引入优雅关闭。业务服务主动断开连接前发送一个关闭帧客户端收到后在 5 秒内完成资源清理并回复服务端收到回复才真正释放连接。这样正常退出的连接不会变成僵尸。2.3 重连退避策略与全局重连风暴防护重连这件事看起来简单——断了就重连。但如果所有客户端都在同一时刻重连那就是事故。一个典型的场景网关发布重启1 万个连接同时断开1 万个客户端同时发起重连请求。网关刚启动还没完全就绪又遇到 1 万并发握手直接就挂了挂了之后另一批客户端也超时重连雪崩。我的防护策略有两层第一层是客户端退避。重连不再是无脑立即重试而是采用退避 抖动的算法第一次 1 秒第二次 2 秒第三次 4 秒……封顶 60 秒每次重连前再加一个 0 到 1000 毫秒的随机抖动。这样即使是同一时间断开的连接也会在 1 到 60 秒的区间内均匀散开避免集中冲击。第二层是服务端过载保护。网关启动后的前几秒如果连接请求量超过预设阈值直接返回服务暂不可用客户端收到这个响应会自动加大退避时间。这相当于给重连风暴加了一个安全阀。3. 消息可靠送达ACK、超时重试与离线补偿推送系统最核心的口碑指标就一个字准。消息要么不推推了就一定要让用户看到。但网络和客户端的不可控因素实在太多推出去和送到完全是两回事。3.1 为什么发出去不等于送到有三种情况会让发出去变成一场空连接假死。服务端以为连接还在实际上消息已经发不出去了数据卡在 TCP 缓冲区里直到超时。客户端崩溃。消息推到了客户端但 App 在处理之前被系统杀掉了。离线场景。用户根本没连上消息根本无处可推。所以可靠送达必须建立在客户端确认之上而不是服务端发送之上。我设计了一套基于 ACK 的确认链路核心数据模型长这样{ mid: msg_20240418_001, from: system, to: user_1024, type: order_update, payload: {}, timestamp: 1713427200000 }每条消息都有一个全局唯一的mid客户端收到消息后需要回一个确认帧{ cmd: ack, mid: msg_20240418_001 }服务端只有收到这个 ACK才认为消息送达成功。3.2 ACK 确认机制的设计ACK 机制说起来简单细节上有个特别重要的问题等待 ACK 的窗口期设多长设太短消息在网络里飞一下就超时了服务端白白重推设太长真丢消息时用户感知延迟会很大。我的做法是设置 10 秒的超时窗口超时后进入重试队列。重试采用指数退避第一次 5 秒后、第二次 10 秒后、第三次 20 秒后……最多重试 3 次。三次都失败就把消息标记为待离线推送触发第 3.3 节的补偿逻辑。这里还有一个容易忽略的细节ACK 不能只做收到消息的确认还得分清到达 ACK和已读 ACK。到达 ACK 是客户端回执证明客户端收到了已读 ACK 是用户点开消息后上报的证明用户看到了。两者混在一起会导致监控数据失真——比如客户端收到但用户没解锁手机到达了却永远不已读如果业务方用已读率衡量推送效果数据就会出现严重偏差。3.3 离线消息拉取与版本号方案用户重连之后服务端如何知道要补哪些消息我用了版本号 拉取的思路简单可靠每个用户维护一个消息版本号version每新增一条消息版本号自增 1客户端本地保存最后收到的版本号last_version重连成功后客户端主动调用一个补偿拉取接口把last_version传给服务端服务端查出 (last_version,current_version] 区间内的所有消息一次性补推这个方案有个很明显的优势不依赖服务端为每个用户维护庞大的离线消息队列消息存到 Redis 或数据库之后按版本号范围查即可。缺点是如果某个用户离线时间很长版本号区间会累积大量消息拉取一次数据量很大。我最后加了个限制超过 100 条且 7 天前的消息合并为一条你有 N 条未读消息的摘要推送避免客户端拉起时被消息洪流淹没。经验提示版本号方案要特别注意并发问题。用户同时登录多台设备手机 平板 网页每台设备都有自己的last_version如果客户端重复推送版本号会出现丢消息或重复消息。我的做法是为每个设备分配独立的device_id版本号按(user_id, device_id)维度记录互不干扰。4. 消息幂等与去重重复推送问题的根因与解法做个推送系统没人想被用户骂怎么又推一遍。但重复推送这个问题几乎每一个做推送的团队都会踩进去。4.1 网络重试导致的消费端重复重复推送的根因往往不在业务侧而在传输链路本身。举一个真实场景服务端把消息推给网关网关把消息写入发送队列但网关还没来得及真正发出去客户端网络闪断了一下发送失败。服务端检测到超时进入重试流程重新推了一遍。此时如果客户端已经收到了第一条消息只是 ACK 没来得及返回那两条一模一样的消息就会同时出现在用户面前。这个问题的本质是在不可靠网络上重试会引入重复不重试会引入丢失。两者不可兼得只能通过去重来兼顾。4.2 基于消息 ID 的幂等方案客户端侧做去重是我第一道防线。每条消息都有一个全局唯一的mid客户端收到消息后先查本地去重表或者 Redis如果mid已经存在直接丢弃不再通知用户。服务端也要做一道幂等校验重试推送前先查一下该消息的目标连接是否已经收到过 ACK。如果收到了说明不需要重推直接标记成功如果没收到先插入一条推送中的记录避免多个重试任务并发重复推送同一条消息。这道校验在 Redis 里实现SET key mid value pushing NX EX 10只有抢到锁的那一个重试任务才真正执行推送。命中重复时我记了一套日志用来统计重复率和网络稳定性。如果某段时间重复率异常升高大概率是客户端网络波动加剧或者是网关超时参数配得太激进。4.3 业务层的最终一致性兜底即使传输层做了去重业务层有时还是会出现重复。比如用户购买商品后同时收到支付成功和订单更新两条通知内容看起来很像但mid不同用户依然会觉得烦。这类问题的解法不在推送系统而在业务系统的消息收敛策略。我给出的建议是在同一业务事件链路上尽量合并通知。支付成功、订单状态变更、物流更新如果三个事件在 5 分钟内都发生推送系统只触发一次聚合通知内容包含所有状态变化。这个策略我是在一次真实反馈后加的——连续 3 条相似推送让用户直接退订了通知后来才意识到推送频控和内容聚合也是体验的关键一环。5. 性能优化与压测从单机千级到万级长连接推送系统跟普通 HTTP 服务有个显著区别长连接占用的资源模式完全不同。普通请求来了就处理、处理完就释放长连接则是一旦建立就长期占用一个文件描述符和一小块内存。所以推送网关的性能瓶颈往往不在 CPU而在文件描述符数量和连接管理的效率上。5.1 网关层优化要点首先调整系统级参数# /etc/sysctl.conf net.ipv4.ip_local_port_range 1024 65535 net.ipv4.tcp_fin_timeout 30 net.ipv4.tcp_tw_reuse 1 fs.file-max 1000000然后是 Nginx 作为 TLS 终结层时的配置worker_processes auto; worker_rlimit_nofile 2000000; events { worker_connections 65535; use epoll; multi_accept on; } http { upstream ws_gateway { server 10.0.0.1:8080; server 10.0.0.2:8080; keepalive 1024; } map $http_upgrade $connection_upgrade { default upgrade; close; } server { listen 443 ssl; http2 on; location /ws { proxy_pass http://ws_gateway; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection $connection_upgrade; proxy_read_timeout 3600s; proxy_send_timeout 3600s; } } }注意一个关键参数proxy_read_timeout。Nginx 默认 60 秒没有数据就会断开代理连接这跟我的心跳策略直接冲突。我把超时调到 3600 秒让心跳来负责真正的存活检测Nginx 只做通道转发不做业务判断。网关服务本身我用的是 Netty核心线程池和业务线程池分离。Netty 的 boss 线程只负责 accept 新连接worker 线程负责 IO 读写业务处理比如消息路由、ACK 处理丢给独立的业务线程池执行。这样即使业务线程池被某个慢操作阻塞IO 线程也不会被拖住连接不至于因为业务阻塞而全部超时。5.2 推送链路的异步化改造初期我的推送链路是同步的业务服务直接调网关接口网关同步拿到连接再同步推送最后同步返回结果。上线后发现一个明显的问题当消息量大时同步链路里任何一环变慢上游生产者都会被拖住数据库连接、线程池很快被打满。改造后我引入了消息队列作为中间缓冲。业务服务把推送任务丢进 RabbitMQ推送消费者从队列里拉取任务再执行路由和下发。异步化之后业务服务和推送网关彻底解耦推送洪峰被队列吸收宁可推送延迟几秒钟也不能让业务服务被推送任务拖死。这里有个数据处理逻辑要注意消费者拉取到任务后要批量处理。单条消息单次发送的效率非常低我改成按用户 ID 聚合同一个用户的连续任务合并成一条批量推送命令减少网络往返。实测这个改动让网关的吞吐提升了将近一半。5.3 压测数据与参数调优压测我用的是开源的压测工具改造后用模拟了 20 万客户端连接重点看三个指标连接成功率、消息推送延迟 P99、网关 CPU 和内存。实测数据如下并发连接数推送速率条/秒P99 延迟网关 CPU网关内存2 万500023ms35%1.2GB5 万1000041ms52%2.8GB10 万1000066ms70%5.5GB20 万10000132ms85%11GB20 万连接时 P99 明显升高排查后发现是 GC 压力上来了——每一条消息都要在内存里创建多个临时对象连接数一多Young GC 频率急剧上升。优化策略是复用发送缓冲区把频繁创建销毁的字节数组对象改为 ThreadLocal 缓存另外把连接对象里不常用的字段比如握手时的一些临时信息延迟初始化减少内存占用。压测还发现了一个很有趣的问题连接数超过 10 万后即使消息量不变网关的 CPU 也会缓慢上涨。分析定位到是心跳扫描线程的算法效率问题。最初心跳检测用 ConcurrentHashMap 遍历所有连接复杂度是 O(n)连接数越多扫描越慢。后来改成时间轮HashedWheelTimer管理心跳超时扫描复杂度降为 O(1)CPU 曲线立刻平了。6. 可观测性建设推送系统到底该怎么监控推送系统最怕什么怕用户说没收到但所有指标都正常。可观测性如果只做到服务没有挂是远远不够的。我上线后逐步积累了下面这一套监控体系基本能做到问题早于用户投诉被发现。6.1 关键指标清单我分四个维度定义监控指标维度指标告警阈值连接层活跃连接数、连接成功率、重连率连接数突降超过 20%传输层推送 P99 延迟、ACK 成功率、超时重试次数P99 超过 200ms消费层队列积压数、消费速率积压超过 1 万条持续 5 分钟业务层到达率、已读率、退订率到达率低于 95%其中 ACK 成功率是我最看重的指标。它衡量的是推出去的消息中有多少被客户端确认收到。如果发现 ACK 成功率下降第一时间排查的应该是客户端网络状态分布而不是服务端——大概率是弱网用户比例升高或者某类手机系统把 App 后台网络给掐了。6.2 一次线上抖动排查全程复盘有一次线上告警推送延迟 P99 从 40ms 飙到 800ms用户侧反馈消息延迟明显。当时我按这个顺序排查先看网关自身指标——CPU、GC、线程池活跃度都正常再看消息队列积压——积压数量并不高说明生产者侧没有堆积看 Redis 慢查询——发现大量 GET 命令的响应时间超过 100ms进一步查 Redis 信息——内存碎片率高得异常原来是路由表的 Key 设置了 3 天 TTL但用户量大大量过期 Key 没有及时清理内存碎片积累导致 Redis 性能劣化修复方案把路由表 TTL 改为心跳周期的 3 倍即 270 秒并开启 Redis 的activedefrag yes。调整后延迟降回正常水平。这个问题的教训是凡是跨组件依赖都要时刻注意对方的资源生命周期是否跟自己的业务模型匹配。一个为短期会话设计的 Key 给了 3 天 TTL就注定会积累垃圾数据。6.3 给后来者的实用建议最后分享几条我在这个项目里用真金白银换来的经验第一一定要做连接数基线管理。给每个网关节点的最大连接数设上限超过上限拒绝新连接并返回稍后重试而不是无底线地硬扛。无底线扩连接最后一定是内存耗竭、整节点宕机。第二消息格式设计要考虑扩展性。我一开始用 JSON 作为消息传输格式后来发现高吞吐场景下 JSON 序列化开销太大。如果量级很大可以考虑改用 Protocol Buffers 之类的二进制格式业务字段加几个不用改协议。第三不要把推送系统的成功标准定义为服务不出错而应该定义为用户真的收到了他想收的消息。技术指标再漂亮用户没感知到都等于零。我在这套系统上踩过的坑十个手指头数不过来。从选型时的纠结到上线前压测的各种调优再到线上真实事故的排查复盘每走一步都对实时推送这四个字理解得更深一些。如果你也在搭建自己的推送系统希望这篇记录能帮你少走几段弯路特别是那些连接生命周期管理、ACK 确认链路的细节真的值得在刚开始设计架构时就认真考虑进去。
返回列表