做消息推送这几年,我踩过不少坑,也积累了一些实打实的经验。今天就把“实时消息推送系统”从选型到落地整个链路拆开聊透,包括技术方案怎么选、架构怎么演进、代码怎么写、上线之后要注意什么。无论你是刚接手的后端开发,还是打算自建推送中台的架构师,这篇文章都值得认真看一遍。
1. 选型背后的思考:为什么我放弃了轮询和长轮询
很多人一上来就问 WebSocket 怎么实现,但实际项目里真正杀人的问题往往在于:你到底需不需要 WebSocket。这件事想不清楚,后面做啥都别扭。
1.1 轮询与长轮询的局限性
我曾经在早期版本里做过一个“伪实时”功能,前端每 5 秒调一次接口拉最新通知,听起来也够用,等用户量到了几万以后就出问题了。每次轮询都产生一次完整 HTTP 请求,加上鉴权逻辑和数据库查询,高峰期服务端 QPS 被拉得很高,而绝大多数响应返回的数据根本没变化,白花花的资源全浪费了。
长轮询比普通轮询好一些:客户端发起请求后,服务端挂住这个连接,等有数据了再返回。这样服务的主动推送时延低很多。可它的问题也很明显:连接挂久了要处理超时,服务端需要额外维护挂起请求的状态,在多实例部署时还需要考虑请求落在哪个节点的问题,复杂度并不低。实测下来的感受是:长轮询属于“能用,但不好扩展”的方案。
1.2 WebSocket、SSE、MQTT 的取舍
消息推送这个场景,主流的实时方案有三个:
- WebSocket:全双工通道,服务端和客户端都能主动发消息,适合聊天、协作编辑、实时通知这类交互型业务。协议基于 TCP,需要自己处理心跳、重连、消息确认这些事。
- SSE(Server-Sent Events):单向通道,只能服务端推给客户端,基于 HTTP,天然支持自动重连,还自带 eventId 断点续传。适合股票行情、工单状态变更这种“服务端主动通知”的场景,实现起来比 WebSocket 轻不少。
- MQTT:基于发布/订阅模型,设计目标是低带宽、弱网环境,适合物联网设备这类场景。如果用户终端的网络不稳定,或者设备能力受限,MQTT 的 QOS 机制会帮你省很多事。
我最终选 WebSocket,原因当中最关键的是业务里不仅有服务端推送,还有用户在线状态、多端互踢这类需要客户端上报消息的需求,SSE 搞不了双向,MQTT 又有点重。WebSocket 在浏览器和移动端的支持都已经很成熟,生态也齐全。假如你只是要单向通知,为了少写代码我其实建议直接用 SSE。
2. 架构设计:单机到集群的演进思路
选型定下来之后,真正的挑战才开始。很多项目第一个版本是单机部署,WebSocket 连接直接放在本地,代码写起来很痛快。但一旦要上多实例,事情立刻就不一样了。
2.1 连接管理怎么做:本地注册表与 Redis 广播
单机环境下,维护所有在线连接最简单的方式就是用一个 ConcurrentHashMap 存 sessionId 到 WebSocketSession 的映射,推送消息的时候遍历这个 Map 挨个发。这方案写起来不到 20 行代码,而且性能极好,本地内存读嘛。
问题出在集群部署后:用户 A 连接在实例 1 上,用户 B 连接在实例 2 上,用户 A 发消息给 B,A 的请求打到了实例 1,实例 1 在自己的本地连接表里根本找不到 B,这消息就发不出去。
解决思路有两条路。一条是引入消息总线,比如 RabbitMQ、Kafka,把“推送给谁、发什么内容”当成一条消息发布到总线,所有实例订阅,谁手里有这个用户连接谁就负责发。另一条是 Redis Pub/Sub + Redis Hash 维护用户连接所在的节点,做定向转发。我的做法是用 Redis Pub/Sub 做广播,理由是接入简单,不依赖额外队列组件,而且推送场景本身对消息不要求持久化和堆积能力,总线只是做一个扇出。
2.2 多实例部署时消息怎么路由
广播方案虽然简单,但如果两个实例同时向同一个用户推送,会出现什么情况?重复消息。这时候你需要一个路由表:用户 ID 和连接所在实例的映射。用户在实例 1 建立连接,就在 Redis 里写一条ws:user:1001 -> instance1;实例 1 收到推送请求,先查路由表,发现用户 1001 不在本地,就把消息转发给实例 1;如果用户多端登录,就要维护一个用户 ID 到多个连接实例的列表。
这个路由表还会带来第二个问题:用户断线时,要记得把路由信息删掉,否则后面推送会一直打到已经不存在连接的实例上,白白浪费一次网络传输。我在实际项目里做过一个妥协:不精确删除,而是让 Redis 里的节点信息带一个过期时间(比如 30 秒),客户端靠心跳续租,服务端靠过期兜底。这样即使删除动作丢了,也不会造成长时间的脏数据。
2.3 消息可靠性的三级保障
在业务上,消息丢失的代价不一样。我习惯把消息可靠性分成三级:通信层可靠性、业务层可靠性和端到端可靠性。
通信层可靠性:WebSocket 协议本身只保证 TCP 层面不出错,但网络中断、服务重启、中间设备超时都会导致消息丢失。这一层要靠心跳检测和重连机制来兜底。
业务层可靠性:服务端在推送之前,先把消息写入数据库(或 Redis),确认客户端真的收到了再标记为已读。我现在做的方案是:入库是一个独立的表,推送动作只负责“尝试发送”,客户端收到消息后返回一个 ACK,服务端收到 ACK 才更新消息状态。
端到端可靠性:这个层级最严格,一般用于支付通知、订单回调这类不能丢也不能重复的场景,除了 ACK 还需要消息幂等去重。客户端收到消息后,先比对本地最新的消息 ID,重复的就不处理。
3. 核心代码实现:基于 Spring Boot 和 WebSocket 的最小可用系统
架构想清楚了,代码才有意义。我拿一个 Spring Boot 的项目举例,从头写一个能跑的 WebSocket 推送系统,包含连接管理、消息推送、离线消息三个核心模块。
3.1 依赖配置与 WebSocket 端点的实现
Spring Boot 接入 WebSocket 非常简单,需要引入一个依赖:
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency>然后定义一个配置类,把 WebSocket 的处理器注册到指定路径上:
@Configuration public class WebSocketConfig implements WebSocketConfigurer { @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(new PushWebSocketHandler(), "/ws") .setAllowedOrigins("*"); } }注意.setAllowedOrigins("*")在正式环境要改成你的前端域名数组,否则跨域问题会一直困扰你,而且安全上也过不去。
3.2 接入认证与连接建立后的生命周期管理
WebSocket 握手协议本质上是一次 HTTP 请求,所以可以利用这个时机做鉴权。我习惯在 URL 上携带一个短时效的 token:ws://youhost/ws?token=xxx。服务端在握手阶段校验 token 的有效性,过期或非法就直接拒绝连接。
然后就是连接管理。我用一个专门的 ConnectionManager 类来管理所有连接,核心数据结构是两层 Map:
@Component public class ConnectionManager { // userId -> 该用户的多个连接 private final ConcurrentHashMap<String, ConcurrentHashMap<String, WebSocketSession>> userSessions = new ConcurrentHashMap<>(); // sessionId -> userId,方便反向查找 private final ConcurrentHashMap<String, String> sessionUserMap = new ConcurrentHashMap<>(); public void addSession(String userId, String sessionId, WebSocketSession session) { userSessions.computeIfAbsent(userId, k -> new ConcurrentHashMap<>()) .put(sessionId, session); sessionUserMap.put(sessionId, userId); } public void removeSession(String sessionId) { String userId = sessionUserMap.remove(sessionId); if (userId != null) { ConcurrentHashMap<String, WebSocketSession> sessions = userSessions.get(userId); if (sessions != null) { sessions.remove(sessionId); if (sessions.isEmpty()) { userSessions.remove(userId); } } } } public List<WebSocketSession> getSessionsByUserId(String userId) { ConcurrentHashMap<String, WebSocketSession> sessions = userSessions.get(userId); return sessions == null ? Collections.emptyList() : new ArrayList<>(sessions.values()); } }为什么存两层 Map 而不是一层?因为同一用户可能多端在线,PC 端、手机端、平板各占一条连接,推送消息时要同时送达所有端。sessionUserMap 是给移除连接时用的,用 sessionId 快速定位 userId,省得遍历。
处理器这边有三个关键回调方法:
public class PushWebSocketHandler extends TextWebSocketHandler { @Autowired private ConnectionManager connectionManager; @Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { // 握手时已经解析好的 userId String userId = (String) session.getAttributes().get("userId"); String sessionId = session.getId(); connectionManager.addSession(userId, sessionId, session); // 推送上线时间、未读消息数量等 } @Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { String payload = message.getPayload(); // 解析消息 type,比如 ACK、心跳、业务上报 } @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { connectionManager.removeSession(session.getId()); // 清理 Redis 路由信息 } }这里有个我踩过的坑:afterConnectionClosed不一定会被及时触发,网络闪断的时候服务端要等 TCP 超时才能感知到。所以光靠这个回调是不够的,必须配合下面的心跳检测。
3.3 消息推送与在线状态管理
推送消息的核心方法大概是这样的:
public void pushToUser(String userId, String messageJson) { List<WebSocketSession> sessions = connectionManager.getSessionsByUserId(userId); for (WebSocketSession session : sessions) { if (session.isOpen()) { synchronized (session) { try { session.sendMessage(new TextMessage(messageJson)); } catch (IOException e) { // 发送失败,记录日志,把 session 关闭 } } } } }为什么要加synchronized (session)?因为 WebSocketSession 不是线程安全的,同一个 session 如果多个线程同时调用 sendMessage,会出现消息交叉错乱的情况。我一开始没加锁,压测的时候发现有概率把两条 JSON 拼在一起发出去,客户端解析直接报错。这个问题在线上排查了很久,后来定位才发现是并发写 session 导致的。
在线状态管理这一块,结合 Redis 比较合适。连接建立后执行:
redisTemplate.opsForValue().set("online:" + userId, sessionId, 30, TimeUnit.SECONDS);客户端每隔 20 秒发一次心跳,服务端收到心跳就重置这个 key 的过期时间。要查用户在不在线,直接查 Redis 就能拿到结果,不用遍历本地连接表,而且这个状态天然可以跨实例共享。
3.4 离线消息与多端登录的处理
如果你只做在线推送,那离线消息必然要接住。我的方案是简洁的消息表:msg_id、user_id、content、status(0 未读 1 已读)、create_time。需要在推送之前先写库,再尝试推送。用户重连后回来拉取一次未读消息列表。
多端登录的问题要复杂些。很多产品希望同一账号只允许一个端在线,新登录的端会把旧端挤下线。实现思路是:新连接建立后,从连接管理器查一下该用户已有的其他 sessionId,向旧 session 发送一个{"type":"KICK"}的消息,然后主动关闭旧连接。如果想要允许同端多端在线(比如两部手机同时登录),上述那个userId + sessionId的映射结构就已经支持了,不需要额外处理。
4. 稳定性设计:心跳、重连、消息确认
WebSocket 看起来简单,真正上生产之后你会明白:连接建立只是一个开始,长连接的稳定性才是大头。这一段内容是我认为整个系统里最值钱的部分。
4.1 心跳机制的实现与参数选择
WebSocket 长连接如果长时间没有数据流动,中间的网络设备(比如 NAT 网关、负载均衡器)会自动回收“空闲连接”,造成服务端和客户端都以为连接还在、实际却已经断开的情况,这就是“幽灵连接”。
心跳的目的有两个:一是及时清死连接,二是让中间设备知道这个连接还活着。我常用的心跳周期是这样的:
- 客户端每 30 秒发送一次
ping消息 - 服务端收到
ping后立即回pong - 服务端如果 90 秒没收到任何消息(不限于 ping),则判定这个连接已死,主动关闭
30 和 90 这个比例是经过权衡的。太频繁了浪费带宽,太稀疏了死连接清理不及时。服务端判断逻辑一般放在一个定时任务里,每隔 30 秒扫描一次所有连接的最后活跃时间,超时就 close。
Spring 的 WebSocket 支持WebSocketSession上的PingMessage和PongMessage,但实际项目中我更建议直接在文本消息里约定心跳字段,因为更容易定位问题,也方便做业务扩展。
4.2 客户端自动重连与幂等去重
客户端 WebSocket 断开是常态,关键看能不能自动恢复。我自己写过前端重连逻辑,经验是必须用“指数退避 + 随机抖动”,即第一次重连等 1 秒,第二次等 2 秒,第三次等 4 秒,最大上限 30 秒,每次再加上一个 0~0.5 秒的随机数。好处是:某个集中报障场景下服务端重启后,所有客户端不会同时发起重连,导致服务端瞬间被打爆。
重连成功后,客户端要主动拉一次增量数据。比如重连之后立刻向服务端请求“从最后一次收到的消息 ID 之后的列表”,把连接断了这段时间漏掉的消息补回来。这个操作要求客户端本地记录一个lastMsgId,每次收到消息后更新。另外客户端收到消息要做去重,因为断线重连和消息重发可能造成同一消息被收到多次。通用做法是维护一个最近收到的消息 ID 集合,消息来了先判断是否处理过,是的话直接跳过。
4.3 消息确认与重发补偿
一个丢消息的高发场景是:服务端调用 sendMessage 成功,但网络延后导致客户端没收到。TCP 层面消息已经发出去了,服务端这里不会报错,但你无法保证客户端真的处理了这条消息。所以要做到“端到端可靠”,ACK 机制避免不了。
具体做法:服务端推送业务消息时带一个msgId,客户端把消息落库或更新本地状态后,回复一条{"type":"ACK","msgId":123}的消息。服务端收到 ACK 后把消息标记为已读,如果 10 秒内没有收到 ACK,就重新发送一次。重试次数不能无限,一般 3 次为上限,超过就标记为“投递失败”,进人工补偿流程。
这个机制引入了一个麻烦:客户端要处理重复消息。所以我说幂等去重和 ACK 是一对孪生兄弟,缺一个另一个就无法正常工作。建议你先把客户端去重实现好,再上 ACK 重试逻辑,否则线上一定会出现“消息重复了”的投诉。
5. 常见问题与排查技巧实录
这一章节的内容全部来自我实际踩坑的现场,不是网上随手能抄到的理论。每个问题都花了很长时间才定位,我整理成速查表形式,方便你遇到类似问题时按图索骥。
5.1 服务器连接数被打满,如何定位瓶颈
一个普遍现象:上线后连接数一涨,CPU 忽高忽低,线程池疯狂报错。这时候不要急着加机器,先看两个指标。
- 文件描述符(file descriptor)占用:每个 TCP 连接都对应一个 fd,Linux 默认 limit 经常是 1024,压测时一会儿就打满了。启动容器或进程前用
ulimit -n 65535调高它。 - 线程数:每个 WebSocket 连接如果独占一个线程,那并发 2000 连接可能就撑不住了。Spring Boot 内置的 Tomcat 默认最大线程数是 200,你可以调整 maxThreads,但我更建议的是换成 Netty 容器,事件驱动模型更适合长连接场景。
我记得有一次线上环境连接数到了 3000 就再也上不去了,排查了一下午,最后发现是tomcat.max-connections的默认值 8192 没问题,反而是 Nginx 的worker_connections只有 1024,客户端都被挡在网关层进不来。所以遇到连接数上不去,先看全链路每一层的连接配置,系统参数、Nginx、容器,逐个排除。
5.2 连接正常但消息推不过去的排查路径
这种问题很诡异,连接建立成功、心跳也正常,但发消息客户端就是收不到。我遇到过的根因大致有三类。
- 消息发到了错误的实例节点上。集群部署时,如果路由信息没更新,消息会推到用户不在线的那个实例。排查方法:看 Redis 路由表里的节点标识是不是指向当前实例,必要时手动删掉让用户重新注册。
- session 已关闭但未被正确移除。网络闪断后,服务端没能及时感知连接关闭,连接管理器里还存着这个 session。清理方法:定期任务遍历所有 session,用
session.isOpen()判断,同时结合心跳最后一次活跃时间,超过阈值强制 remove。 - 序列化或消息体格式问题。客户端解析不了服务端发的 JSON,会静默地丢弃消息或者抛异常但不退出。这个看似低级,实际最容易发生。建议服务端记录推送的消息体日志,客户端在 onmessage 里加一个全局错误捕获,定位到这一步会快很多。
5.3 Nginx 配置与容器部署的注意事项
用 Nginx 做反向代理时,WebSocket 有一项特殊配置必须加上,否则连接建立后几秒就会断开,这是最常见的新手问题:
server { listen 80; server_name push.example.com; location /ws { proxy_pass http://backend_api; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_read_timeout 3600s; proxy_send_timeout 3600s; } }proxy_read_timeout是关键参数,默认 60 秒。WebSocket 连接建立后如果 60 秒内没有数据交互,Nginx 就会掐断它。你没有心跳机制的话,这条断连会表现为“客户端时不时掉线,重连后又好”。加心跳能缓解,但建议还是把超时时间调大。
部署容器实例时还要注意优雅停机。Kubernetes 滚动更新把旧 Pod 杀掉的时候,如果直接 SIGKILL,客户端端会表现为“突然断线”。正确处理是配一个preStop钩子,让旧实例延迟几秒再退出,给客户端一点时间走重连逻辑,同时服务端等待 5 秒处理完当前连接上的消息。我用的是 Spring 的ContextClosedEvent,在里面关闭所有 session 并清理 Redis 路由。
5.4 消息乱序与并发写 session 的坑
消息乱序这个问题,很多人会忽略。WebSocket 本身是保证消息顺序的,但如果你在多线程环境里并发发送消息,就可能出现后发先至。我在 3.3 小节提过synchronized (session)的解法,这里再补充一个:如果对顺序有更强要求,比如聊天类业务,建议做成串行队列。每个 session 关联一个单线程的 Executor,所有发送任务进队列,保证同一条连接上的消息严格有序。代价是内存会多一点,但换来的是顺序性可靠。
我同事用的是一个更轻的做法:发送消息前给每条消息加一个序号,客户端收到后如果序号不连续,就主动请求补发。这个方案好在服务端不用维护复杂的队列,坏处是客户端逻辑复杂一些。看你的业务要不要这么严谨,非强实时业务用锁就够了。
6. 项目上线后还需要做的几件事
写完代码、联调通过不代表项目结束。消息推送这种长连接业务,上线后要盯的量和工作方式都有讲究。
6.1 监控指标与告警策略
先说一下我最终沉淀下来的核心监控指标:
- 当前活跃连接数:异常波动(瞬间暴跌或暴涨)都要告警
- 连接建立速率:秒级新建连接的速率,反映客户端重连风暴
- 消息推送吞吐量:每秒推送的消息条数,和业务量挂钩
- 消息投递成功率:推送成功数 / 应推送总数,低于 99% 就要查了
- 心跳超时率:单位时间内心跳超时的连接占比,反映网络质量
告警阈值不需要定得太严格,我用的经验是:连接数下降超过 30% 触发高优先级告警;消息投递成功率连续 3 分钟低于 95% 触发告警;心跳超时率持续 5 分钟超过 10% 触发告警。太灵敏反而会频繁打断你的工作节奏。
6.2 压测方法和容量评估
压测 WebSocket 系统不能只测 HTTP 接口,要专门用支持 WebSocket 的压测工具,我用的是 JMeter 的 WebSocket Sampler 插件和 Gatling。测试场景至少覆盖三块:单实例能支撑多少连接数、连接并发建立时 CPU 的表现、消息广播时的吞吐量和延迟。
最后给你一个粗略的容量参考:一个 4C8G 的实例,用 Netty 容器,支撑 3~5 万条长连接问题不大。瓶颈通常在内存,每个 WebSocketSession 大概占 2~5KB,5 万连接就是 100~250MB 内存,还要留足给业务逻辑。如果你用 Tomcat 容器,同样的配置可能要打七折。
6.3 后续拓展:从自研组件到接入消息推送中台
如果你的业务量再往上走,自研这套东西的边际成本会越来越大。连接集群、跨机房容灾、推送链路跟踪、推送效果分析,每一块都是工作量。到那个阶段,我建议你把系统做厚:底层用消息队列做削峰填谷,中间沉淀一层“推送编排”能力,上层对业务提供统一 API,支持普通推送、批量推送、定时推送,同时把推送状态回调给业务方。从这个角度说,一开始设计的时候就不要把业务逻辑焊死在 WebSocket 连接上,尽量通过消息体里的 type 字段做解耦,不然后面每次加一类推送需求都要发一次版。
这套系统从最初的轮询到后面的 WebSocket 集群方案,中间经历了两次比较大的重构,每次都是被真实流量逼出来的。我个人最大的体会是:技术上没有银弹,实时推送的方案选型必须服从业务场景。如果你的业务是低延迟强交互,WebSocket 值得投入;如果只是单向通知,SSE 足够;如果网络环境很差,那就老老实实研究 MQTT 的质量等级设置。最后再分享一个小技巧:上线初期把消息体日志全部打开,包括连接建立、连接关闭、推送成功、推送失败,持续观测两周。虽然日志量大一点,但你会对整个系统的脾气摸得清清楚楚,后面再排查问题会快好几倍。