消息推送系统的性能优化实战:从单机百万连接到分布式推送中台的架构演进

发布时间:2026/7/22 0:23:52

消息推送系统的性能优化实战:从单机百万连接到分布式推送中台的架构演进 消息推送系统的性能优化实战从单机百万连接到分布式推送中台的架构演进一、当推送带宽触及物理极限单机百万连接的不可逾越之墙消息推送系统的核心挑战可以用一个数字概括一个 Linux 服务器默认的端口范围是 65535当你要同时维持与 100 万客户端的 TCP 长连接时传统的 BIO 模型一个连接一个线程首先在内存上就吃不消——每个线程栈默认 1MB100 万线程就是约 1TB 内存。即便用 NIO 解决了线程问题单机的网络带宽、CPU 中断处理能力、文件描述符数量都有硬上限。某社交平台的消息推送系统在 DAU 从 500 万增长到 3000 万的过程中经历了三次架构升级。第一代架构基于 Netty 的单机方案单节点能扛住约 20 万长连接受限于 4C8G 实例的内存带宽第二代引入了集群化部署和一致性哈希路由将连接均匀分散到 20 台节点上第三代则演进为分布式推送中台将连接管理、消息路由、推送策略、下行加速完全解耦。本文聚焦第二次到第三次升级过程中的关键技术决策大规模长连接管理策略、消息的可靠可达性保障、以及推送策略的智能降级机制。二、推送中台的三层解耦架构连接层、路由层、策略层的独立伸缩升级后的推送中台将系统拆分为三个可以独立伸缩的层次连接层是系统的前线每个节点使用 Netty 的 EpollEventLoopGroup 来管理数十万 TCP 长连接。连接层的核心任务是维持连接的存活状态不关心消息内容。每个连接在建立时生成一个全局唯一的 Session ID由路由层负责 Session ID 到连接节点的映射。路由层是系统的调度中心。当业务方发送一条推送消息时消息网关将消息投递到 RocketMQ 中。推送调度器消费消息后根据推送类型选择路由策略广播消息直接发送到所有连接节点单播消息通过一致性哈希定位目标用户的 Session 所在节点标签推送则先查询 Redis 中的标签索引获取用户列表再按节点分组批量发送。策略层是系统的调节阀。它负责三项核心策略下行速率控制——根据每个连接节点的 CPU 和网络带宽动态限制每秒推送量防止突发流量打崩连接层心跳自适应——根据连接的活跃度动态调整心跳间隔活跃连接 30 秒一次静默连接延长到 3 分钟一次降低电池消耗和服务端压力离线消息存储——当用户不在线时将消息暂存到 Redis上线后批量拉取。三、Netty 长连接管理与可靠送达的核心实现以下是连接层中 Netty 长连接管理和消息确认机制的核心代码/** * 推送连接管理器 * * 核心职责 * 1. 管理数十万级别长连接的生命周期 * 2. 维护 Session - Channel 的映射关系 * 3. 实现消息的 ACK 确认与重试机制 */ Component public class PushConnectionManager { /** * 使用 ConcurrentHashMap 维护在线连接映射 * Key: userId_appId_deviceId 组成唯一的 SessionId * Value: Netty Channel 引用 */ private final ConcurrentHashMapString, Channel sessionChannelMap new ConcurrentHashMap(); /** * 反向索引Channel - SessionId用于连接断开时快速清理 */ private final ConcurrentHashMapChannelId, String channelSessionMap new ConcurrentHashMap(); /** * 待确认消息队列msgId - (SessionId, 消息内容, 重试次数, 发送时间) */ private final ConcurrentHashMapString, PendingMessage pendingMessages new ConcurrentHashMap(); // 重试配置 private static final int MAX_RETRIES 3; private static final long RETRY_INTERVAL_MS 3000; PostConstruct public void startAckChecker() { // 启动后台 ACK 检查线程 ScheduledExecutorService checker Executors.newSingleThreadScheduledExecutor( r - new Thread(r, ack-checker)); checker.scheduleAtFixedRate(this::checkPendingMessages, 5, 3, TimeUnit.SECONDS); } /** * 发送消息并注册 ACK 监听 */ public SendResult sendMessage(String sessionId, String payload) { Channel channel sessionChannelMap.get(sessionId); if (channel null || !channel.isActive()) { return SendResult.offline(用户不在线); } String msgId UUID.randomUUID().toString().replace(-, ); PushMessage message PushMessage.builder() .msgId(msgId) .payload(payload) .timestamp(System.currentTimeMillis()) .build(); try { ChannelFuture future channel.writeAndFlush( new TextWebSocketFrame(JSON.toJSONString(message))); // 注册待确认消息 pendingMessages.put(msgId, PendingMessage.builder() .sessionId(sessionId) .message(message) .retryCount(0) .firstSendTime(System.currentTimeMillis()) .build()); // 异步监听发送结果 future.addListener((ChannelFutureListener) f - { if (!f.isSuccess()) { log.warn(消息发送失败, msgId{}, sessionId{}, msgId, sessionId); handleSendFailure(msgId, sessionId); } }); return SendResult.success(msgId); } catch (Exception e) { log.error(消息发送异常, sessionId{}, sessionId, e); pendingMessages.remove(msgId); return SendResult.failed(发送异常: e.getMessage()); } } /** * 处理客户端 ACK 确认 */ public void handleAck(String msgId, String sessionId) { PendingMessage pending pendingMessages.remove(msgId); if (pending ! null) { log.debug(消息确认送达, msgId{}, 耗时{}ms, msgId, System.currentTimeMillis() - pending.getFirstSendTime()); } } /** * 定时检查未确认消息执行重试或标记失败 */ private void checkPendingMessages() { long now System.currentTimeMillis(); ListString toRemove new ArrayList(); for (Map.EntryString, PendingMessage entry : pendingMessages.entrySet()) { PendingMessage pending entry.getValue(); if (pending.getRetryCount() MAX_RETRIES) { // 超过最大重试次数标记为投递失败触发离线存储 toRemove.add(entry.getKey()); storeOfflineMessage(pending); continue; } if (now - pending.getFirstSendTime() RETRY_INTERVAL_MS * (pending.getRetryCount() 1)) { // 执行重试 Channel channel sessionChannelMap.get(pending.getSessionId()); if (channel ! null channel.isActive()) { pending.incRetryCount(); channel.writeAndFlush(new TextWebSocketFrame( JSON.toJSONString(pending.getMessage()))); } else { // 连接已断开消息转离线存储 toRemove.add(entry.getKey()); storeOfflineMessage(pending); } } } toRemove.forEach(pendingMessages::remove); } }ACK 确认机制是可靠推送的基石。我们采用的策略是每条消息发出去后注册到 pending 表等待客户端回 ACK如果 3 秒内未收到 ACK 则重试最多重试 3 次3 次后若仍失败消息转入离线存储。这里的关键参数——3 秒重试间隔和 3 次上限——是通过生产环境的延迟分布数据确定的。在我们的场景中P99 的端到端延迟约 800ms大部分消息在 1.5 秒内完成送达和确认3 秒的间隔给予了足够的容错又不会产生过多延迟。四、弱网环境下的推送可靠性方案的适用边界这套方案在弱网环境下面临明显的考验。当一个用户在高铁上反复经历 4G/5G 切换时TCP 连接会频繁断开重建。每次重建都触发离线消息的批量拉取可能引发消息风暴——同时拉取数百条离线消息瞬间占满客户端的蜂窝带宽。解决策略有两个方向。一是引入消息分层机制将消息分为 P0必须送达如订单状态变更、P1期望送达如系统通知、P2尽力送达如营销推送。P1 和 P2 消息在离线超过指定时长后直接丢弃减少离线消息的堆积量。二是做离线消息的分页拉取——每次只拉取最近 20 条用户向上滑动时再加载更多避免一次性推送全部离线消息。另一个需要警惕的场景是大 V 推送风暴。当一个拥有千万粉丝的账号发了一条动态需要瞬间推送到千万台设备的消息投递量可能在数分钟内达到百万 QPS。即使使用消息队列削峰连接层的 Channel 写入也会成为瓶颈。我们的应对策略是引入推送优先级和速率限制——高优先级消息如 IM 私信全量推送低优先级消息如动态更新提醒按连接节点的处理能力做发送速率限流峰值时段可接受数分钟级别的到达延迟。在成本方面长连接方案的一个隐藏成本是服务器的带宽费用。以单节点 20 万连接、每条心跳消息 50 字节、心跳间隔 30 秒为例心跳流量约为 20万 × 50B × (1/30s) × 2(上下行) ≈ 667KB/s。看起来不高但如果心跳做加密增加 30% 体积、加上业务消息的下行流量一个 100Mbps 带宽的节点连接数上限大约在 15 万左右。带宽成为比 CPU 和内存更早到达瓶颈的资源。五、总结消息推送系统的性能优化不是靠单一技术突破的而是在连接管理、消息路由、推送策略三个层面做系统性设计的结果。Netty 的 Epoll 模型解决了并发连接的基础能力一致性哈希路由解决了集群化的连接分发问题ACK 确认 离线存储解决了消息的可靠可达性而下行速率控制则在突发流量下保护了连接层的稳定性。落地建议第一阶段先验证 Netty 单机容量找到当前硬件配置下的连接数上限和带宽天花板建立容量模型第二阶段引入路由层和 RocketMQ 做消息削峰实现连接层和策略层的解耦第三阶段补齐离线消息存储、消息优先级分级和消息的生命周期管理。核心监控指标包括在线连接数、消息送达成功率ACK 率、P99 推送延迟、心跳异常率以及连接节点的 CPU/内存/带宽使用率。

相关新闻