首页 / 资讯中心 / 文章详情

WebSocket集群消息不丢:Spring Boot集成RabbitMQ广播实践

WebSocket集群消息不丢:Spring Boot集成RabbitMQ广播实践 ★ FEATURED ARTICLE
1. 从单体到集群WebSocket推送为什么会丢消息先说一个我实际踩过的坑项目早期是一个单体Spring Boot应用用的spring-boot-starter-websocket加STOMP前端连上来就完事服务端要推消息直接SimpMessagingTemplate.convertAndSend(/topic/xxx)一切岁月静好。后来业务量上来单机扛不住我们把服务扩成三节点前面挂Nginx做负载均衡——噩梦就从这里开始了。为什么会丢消息核心原因在于WebSocket是一个长连接它天生粘人。客户端一旦完成STOMP握手连接就固定在某个具体的服务实例上这个实例的会话表里保存了这条Session。当你往/topic/xxx发消息时消息只会经过你连接的那个Broker转发给本地会话——如果客户端A连的是实例1而实例2上的某个业务逻辑触发了推送消息发出去就是石沉大海因为实例2的会话表里根本没有客户端A的Session。这就像你在一个小卖部里喊人老板认识所有常客嗓门一大人人都听得见。但当你把生意做到三个分店总店喊一嗓子分店的客人自然听不到。单体架构下Spring的简单消息代理Simple Broker是进程内存级别的它不知道、也管不到其他进程的会话。解决思路一般有三条共享Session存储 自研路由把Session信息放到Redis推送时找到Session所在实例再发起内部调用。听起来合理但实现起来要处理跨节点通信、Session失效同步、粘性路由失败重试复杂度不低。粘性会话Sticky Session让同一个客户端永远打到同一个实例。这在短连接场景没问题但WebSocket是长连接一旦某个实例重启、发版滚动更新粘在上面的连接全部断线客户端必须重连集群的高可用打了折扣。外部消息代理 广播订阅用RabbitMQ或ActiveMQ作为STOMP Broker所有实例订阅同一个Topic谁收到消息谁负责把消息转给本地会话。客户端任意连到哪个实例都能收到推送。我最终选的是第三条。理由很直接它在不改动客户端连接方式的前提下解决了一致性问题而且Spring Boot对RabbitMQ STOMP的支持已经非常成熟不需要自己造轮子。下面我把整条链路拆开讲。注意这个方案不是客户端到客户端的聊天室方案它解决的是服务端主动单向推送的消息可达性问题——比如订单状态变更、系统通知、告警触发这类场景。如果是实时双向通信方案还得调整。2. 集群下消息路由的原子性为什么每个实例都收到不丢消息我们先澄清一个容易混淆的概念。很多人听到广播就担心同一个消息发给所有实例每个实例都往客户端推一遍客户端会不会收到多条重复消息答案是不会。因为每个实例只转发属于它本地会话的消息。STOMP的Topic模型在服务端看来订阅者是会话Session而不是客户端。实例1收到一条推送到/topic/alerts的消息它遍历自己的Session表发现有3个会话订阅了这个Topic就推给这3个连接实例2上也订阅了同样的Topic但它本地只有2个会话就推给这2个。这5个会话对应5个不同的客户端消息各推一次互不重复。这里的关键在于消息的原子性——消息必须被集群中的每一个实例都接收到再有各自的本地会话去消化。只要有一个实例没收到消息它上面的会话就会漏掉这条推送。所以集群推送的核心不是发给谁而是**如何保证所有实例都收到**。这时候就需要一个外部消息代理来承担扇出Fan Out的职责。以RabbitMQ为例它的Topic Exchange主题交换机天然支持通配符路由——你定义一个队列绑定关系让每个Spring Boot实例都创建一个同名队列并绑定到同一个Exchange上发送端只要往这个Exchange发一条消息RabbitMQ会把它复制到所有绑定的队列中每个实例的监听器自然都能消费到。用生活化的比喻RabbitMQ就是那个群发邮件的邮局每个服务实例是邮局分店客户是收件人。邮局只要把同一封信投递到每家分店的邮箱里分店再把信送到各自片区的客户手中就能保证所有客户都收到信。所以这一阶段的架构变成这样三层结构客户端 → Spring Boot实例含STOMP端点 本地Broker转发 → RabbitMQ BrokerSpring Boot实例既是STOMP服务端也是RabbitMQ的客户端发送端不再直接往本地Broker推消息而是通过RabbitMQ的Topic Exchange发布每个实例监听同一个队列拿到消息后用本地的SimpMessagingTemplate往本地会话推送这里有一个重要的点Spring Boot的STOMP配置里你可以不再依赖本地的Simple Broker而是指定用RabbitMQ作为Broker。也就是说客户端发来的/topic订阅请求会由RabbitMQ去管理订阅关系而不是实例本地管理。这样一来客户端A连的是实例1但它订阅/topic/alerts这个动作会被同步到RabbitMQ——实例2推送消息时RabbitMQ知道实例1上有客户端A的订阅就会把消息投递给实例1实例1再转给客户端A。本质上你连订阅关系都被集群化了消息丢失的可能性进一步降低。不过这里要注意全走RabbitMQ Broker会引入额外的网络跳数和延迟对于单向通知类业务场景如告警、站内信本地Simple Broker 外部MQ广播的混合模式反而更轻量、更容易排查问题。我这次方案选的就是混合模式客户端STOMP端点由Spring Boot提供订阅关系由本地的Simple Broker管理跨实例的广播走RabbitMQ的Topic Exchange。这个思路的好处是客户端不需要感知RabbitMQ的存在服务端推送链路边界清晰出问题时只需看Producer是否把消息发到了Exchange以及Consumer是否从Queue里拿到了消息这两段即可排查成本低。3. 手写Token握手校验为什么不能用Spring Security的默认过滤器STOMP握手本质是一次HTTP请求但握手成功后连接就升级成WebSocket长连接后续帧不再走HTTP协议。这意味着——你没法用Spring MVC里的HandlerInterceptor或Filter去拦截STOMP帧来做认证因为那些都是HTTP层的东西。我在最初做Token认证时踩过一个坑用了Spring Security在SecurityFilterChain里配置了/ws/**放行又在WebSocket握手拦截器里做Token校验。看起来逻辑通顺但等到联调时发现某些客户端的Token是在握手后的第一条CONNECT帧里带过来的握手阶段根本拿不到。这样一来握手阶段拦截器校验失败连接直接被拒前端一脸懵。标准做法是不要依赖握手的HandshakeInterceptor做核心认证而是用ChannelInterceptor拦截CONNECT帧在该帧里校验Token。为什么因为STOMP协议规定CONNECT帧是客户端在WebSocket建立之后、正式建立STOMP会话之前发送的第一帧里面可以携带Authorization头如果你用的是自定义Header也可以放在nativeHeaders里。这个时机正好是认证的黄金窗口——连接已经建立但尚未进入业务消息收发阶段此时拒绝连接是干净且无副作用的。认证链路大致是客户端在WebSocket URL上带Token参数如ws://host/ws?tokenxxx——方便握手阶段做快速预检握手拦截器拿到Token后只做格式预检是否为空、是否过期不查库避免握手链路过重客户端发送CONNECT帧携带Authorization: Bearer tokenChannelInterceptor实现类拦截CONNECT帧从StompHeaderAccessor里取出Auth头解析Token校验签名和有效期加载用户信息认证通过后把用户信息写入StompHeaderAccessor的sessionAttributes或user属性后续业务消息里直接取用认证失败则返回ERROR帧并关闭连接这里有个细节要强调ChannelInterceptor拦截的是整个STOMP通道上的所有帧不仅仅是CONNECT。所以你在实现时要判断StompCommand.CONNECT只在这个命令上做认证逻辑否则后续每一条消息都走一遍Token校验浪费不说还可能误伤订阅和发送。写完拦截器后还需要在WebSocket配置里把它注册为clientInboundChannel的拦截器Configuration EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { private final TokenChannelInterceptor tokenChannelInterceptor; public WebSocketConfig(TokenChannelInterceptor tokenChannelInterceptor) { this.tokenChannelInterceptor tokenChannelInterceptor; } Override public void configureClientInboundChannel(ChannelRegistration registration) { registration.interceptors(tokenChannelInterceptor); } Override public void registerStompEndpoints(StompEndpointRegistry registry) { registry.addEndpoint(/ws) .setAllowedOriginPatterns(*) .withSockJS(); } Override public void configureMessageBroker(MessageBrokerRegistry registry) { // 客户端订阅前缀 registry.enableSimpleBroker(/topic, /queue); // 服务端消息前缀 registry.setApplicationDestinationPrefixes(/app); } }这个配置里有个值得注意的设计enableSimpleBroker(/topic, /queue)——/topic用于广播推送/queue用于点对点推送。在集群场景下/queue一定要小心用因为Simple Broker的点对点也是进程内存级别的客户端A在实例1上订阅了一个/queue/xxx实例2往/queue/xxx推消息实例1是收不到的。如果需要跨实例点对点应该改用RabbitMQ的Queue模式或者把点对点也设计成广播每个实例只关心自己是否有匹配的本地会话其实配合Session级用户映射反而更简单。4. RabbitMQ的接入与自动配置一套配置解决订阅关系和消息扇出我选择RabbitMQ作为外部Broker有一个很实际的考量Spring Boot的spring-boot-starter-amqp对RabbitMQ的封装足够成熟而且Spring官方的spring-messaging模块里SimpleMessageBroker和RabbitMessageBroker的切换成本非常低。但要提醒的是直接用RabbitMQ作为STOMP Broker和混合模式下的消息扇出中间件是两回事。前者需要启用RabbitMQ的STOMP插件客户端直接把STOMP帧发给RabbitMQ的61613端口后者客户端仍然把STOMP帧发给Spring Boot服务Spring Boot再作为生产者和消费者与RabbitMQ打交道。我做的是后者理由之前说过客户端只需要知道一个入口地址Nginx代理后的WebSocket端点不需要关心背后有几台RabbitMQ、哪个队列在哪灵活性和可维护性都好很多。具体实现时需要在项目里引入依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-websocket/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-security/artifactId /dependency dependency groupIdio.jsonwebtoken/groupId artifactIdjjwt-api/artifactId version0.11.5/version /dependency然后配置文件里加上RabbitMQ的连接信息spring: rabbitmq: host: 10.0.0.12 port: 5672 username: admin password: admin123 listener: simple: acknowledge-mode: manual concurrency: 3 max-concurrency: 8这里我把acknowledge-mode设成了manual——手动确认。为什么因为WebSocket推送是异步的消息到达监听器后推送给客户端的过程可能很快但你不能假设它一定成功。如果监听器方法抛出异常自动确认模式下消息会被丢弃推送就丢了手动确认模式下你可以先尝试推送推送失败就执行basicNack让消息重回队列由另一个线程重试。RabbitMQ的消息拓扑我这样设计的交换机push.exchange类型是topic队列push.queue.v1绑定键是push.*发送端convertAndSend(push.exchange, push.notice, payload)每个Spring Boot实例都声明同一个队列push.queue.v1绑定同一个交换机为什么要同一个队列而不是每个实例一个队列这涉及到RabbitMQ的一个特性同一个队列被多个消费者监听时消息是负载均衡分发的Work Queue模式而不是广播。在推送场景下我们不希望这样——我们希望每个实例都能收到消息所以必须每个实例声明一个独立队列或者让每个实例声明同一个主题队列的副本。正确做法是每个实例声明一个带唯一后缀的队列比如push.queue.instance1、push.queue.instance2都绑定到push.exchange上绑定键都用push.*。这样发送端发一条消息RabbitMQ会根据绑定关系复制到每个实例的队列里实现广播效果。Bean public Queue pushQueue() { // queueName 是每个实例启动时自己生成的前缀例如 push.queue. UUID.randomUUID() return new Queue(queueName, true, false, true, Map.of(x-ha-policy, all)); } Bean public TopicExchange pushExchange() { return new TopicExchange(push.exchange, true, false); } Bean public Binding pushBinding() { return BindingBuilder.bind(pushQueue()).to(pushExchange()).with(push.*); }消费者这一侧用RabbitListener监听唯一队列收到消息后调用SimpMessagingTemplate.convertAndSend把消息推给本地订阅者RabbitListener(queues #{pushQueue.name}) public void onPushMessage(Message message, Channel channel) throws Exception { try { PushPayload payload objectMapper.readValue(message.getBody(), PushPayload.class); // 推给订阅了对应主题的所有本地会话 messagingTemplate.convertAndSend(/topic/ payload.getTopic(), payload.getData()); // 手动确认 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); } catch (Exception e) { log.error(push message error, e); channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); } }注意这里messagingTemplate.convertAndSend的目标是/topic/...这个主题是客户端直接订阅的主题而RabbitMQ那边的绑定键push.*只是内部路由用的。层与层之间是解耦的消息在进入Spring Boot之后由本地Broker决定往哪些会话推送。5. 集群环境下的负载均衡与粘性会话该配的还是要配有人会问既然广播机制已经保证了消息不丢那Nginx还需要粘性会话吗答案是建议配置但不能依赖。原因有两个第一个原因是WebSocket连接升级的流程。客户端先发起GET /ws的HTTP握手请求Nginx需要把这条请求转发给某个后端实例握手成功后升级成WebSocket长连接。如果Nginx没有配置Upgrade和Connection头这个请求会被当成普通HTTP请求处理WebSocket握手直接失败。这是最低配必须做。第二个原因和握手阶段带Token有关。如果Token是放在URL参数里的那么每次重连都是无状态的路由到哪个实例都一样不需要粘性会话。但如果你在CONNECT帧里带Token而Nginx做的是HTTP层面的负载均衡它根本看不到STOMP帧里的内容——这和分析TCP payload一样不现实。所以粘性会话更多是为了减少WebSocket频繁重连带来的握手开销而不是解决消息丢失问题。Nginx的配置片段upstream ws_backend { server 10.0.0.1:8080; server 10.0.0.2:8080; server 10.0.0.3:8080; ip_hash; } server { listen 80; location /ws { proxy_pass http://ws_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_set_header Host $host; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; proxy_read_timeout 3600s; proxy_send_timeout 3600s; } }ip_hash在这里承担了粘性会话的职责——同一个客户端IP的请求总是转发到同一个后端实例。但要注意如果客户端走的是移动网络IP可能会频繁变化或者多个用户共享同一个出口IPip_hash就会不够精准。更高级的做法是sticky模块基于Cookie做会话保持不过对于WebSocket场景Cookie方案略显鸡肋——毕竟握手之后连接是长连接只要不重启一直挂在那个实例上就够了。真正需要处理的是实例重启导致的连接断开。游戏规则是这样的你无法保证实例不重启所以客户端必须实现自动重连机制。前端在onclose事件里做指数退避重连重连时会带上同一个Token重新握手。这时粘性会话的作用就体现出来了——如果重连的请求被路由到了另一个实例消息推送也不受影响因为广播机制保证每个实例都能消费到RabbitMQ里的消息只是握手成本高了一点点。所以严格来说有粘性会话性能好一点没有粘性会话功能也不会坏。这就是我说的建议配置但不能依赖。6. Token认证与clientInboundChannel别再走HTTP Filter的老路这一节我想展开讲讲Token认证的实现细节因为这是容易出看似对实则错的地方。先纠正一个常见误解有人会在WebSocketConfig里通过addInterceptors注册HandshakeInterceptor然后在beforeHandshake方法里做Token解析。这个方法是能拿到Token的但它有个致命问题——它只处理握手阶段的HTTP请求参数。如果你约定Token必须放在Authorization头里那么常规的JavaScript WebSocket API并不能自定义Header你只能放在URL参数或protocol里。前端代码稍不注意就把这个环节绕过去了。更稳妥的方式是三阶段校验握手阶段可选但推荐从URL参数里取token只检查是否为空、格式是否合法做一个低成本的快速预检。CONNECT帧阶段必须从StompHeaderAccessor的nativeHeaders里取Authorization头做完整校验签名有效期用户信息加载。订阅/发送阶段按需针对SUBSCRIBE或SEND命令做细粒度的权限控制。第二阶段的实现大致是这样public class TokenChannelInterceptor implements ChannelInterceptor { private final JwtTokenProvider tokenProvider; private final UserDetailsService userDetailsService; Override public Message? preSend(Message? message, MessageChannel channel) { StompHeaderAccessor accessor MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); if (accessor ! null StompCommand.CONNECT.equals(accessor.getCommand())) { String token extractToken(accessor); if (token null || !tokenProvider.validateToken(token)) { throw new AuthenticationCredentialsNotFoundException(invalid token); } String username tokenProvider.getUsername(token); UserDetails userDetails userDetailsService.loadUserByUsername(username); // 关键是这一行把用户信息放进 accessor 中后续可以通过 accessor.getUser() 获取 accessor.setUser(userDetails); } return message; } private String extractToken(StompHeaderAccessor accessor) { ListString authHeaders accessor.getNativeHeader(Authorization); if (authHeaders null || authHeaders.isEmpty()) { return null; } String header authHeaders.get(0); if (header.startsWith(Bearer )) { return header.substring(7); } return null; } }这里的关键点是accessor.setUser(userDetails)。设置之后后续的业务代码里通过SimpMessagingTemplate.convertAndSendToUser(username, /queue/private, payload)可以实现点对点推送——它的底层原理是把用户名和Session ID关联起来在Simple Broker的本地会话表里做映射。单节点下一切正常集群下convertAndSendToUser是否能跨节点送达取决于用户是否恰好连接在本实例上。所以如果需要跨节点的点对点推送要么用RabbitMQ作为Broker要么把点对点改成广播给所有实例每个实例根据本地Session表判断是否推送。我在实际项目中更倾向于后者逻辑更直接。另外值得注意的一点Token过期后连接并不会自动断开。因为STOMP长连接一旦建立后续帧不再做HTTP层面的认证只要你不实现定期校验Token的逻辑过期的Token依然可以继续收发消息。我在生产环境里加了一个心跳机制客户端每30秒发送一个/app/heartbeat消息服务端在ChannelInterceptor里拦截这个命令并检查Token是否过期过期就返回ERROR帧强制断开连接。这个机制对那些长期挂机不关闭页面的用户特别重要。7. 实例重启与连接补偿状态恢复思路集群滚动发布或者实例崩溃时长连接一定会断开这是WebSocket的天然宿命。我们能做的不是避免断线而是保证断线恢复后的消息不丢。这里有个矛盾如果消息在实例重启期间推送那一刻没有实例持有客户端的Session消息即使被RabbitMQ广播到所有实例也找不到接收者。简单粗暴的解决方式是宁可多不可少——推送方在发出消息时把消息持久化到Redis里带上一个全局唯一的消息ID客户端重连成功后上报自己最后收到的消息ID服务端把差值补发。这个机制我称之为时间窗口补偿。实现思路不复杂推送方发送消息前先写入Redis Streamkey按业务维度区分value包含消息体、发送时间、消息ID。RabbitMQ负责把消息广播到各实例各实例推送给本地会话。客户端在每次推送消息里带上msgId客户端存储最近收到的msgId重连后的第一条消息带上它。服务端提供一个/app/pullMissing?msgIdxxx的端点根据客户端上报的消息ID从Redis Stream里读取之后的消息通过convertAndSendToUser补发。这个方案的复杂度确实上去了但对可靠推送有硬性要求的业务比如交易结果通知、工单状态更新这一步是绕不过去的。如果只是普通站内信、状态提醒你可以接受极端情况下的少量丢失那把这个机制做成可选的补偿开关就好不必常态启用。我个人的实现经验是先把基础推送链路做稳RabbitMQ广播 手动ACK再把补偿机制作为第二优先级。因为补偿机制本身依赖消息已经持久化这个前提而如果广播链路本身就丢消息补偿机制反而会掩盖问题让你排查起来更费劲。递进式的架构演进比一步到位更稳。8. 生产环境实测与踩坑总结几个让我排查到凌晨的细节最后分享几个在生产环境里踩过的坑。每一个都让我和团队吃过苦头写出来帮大家提前绕开。第一个坑RabbitMQ广播变成工作队列。前面提到每个实例要声明独立队列但如果你的实例是在PostConstruct里动态创建队列的并且队列名写死了比如push.queue那么三个实例声明的是同一个队列名RabbitMQ会认为这是同一个队列的三个消费者——消息只会被其中一个消费者收到。表面上看消息好像没丢但实际推送覆盖率急剧下降。排查这个问题时我先看了RabbitMQ管理台发现三个实例的Connection都连着同一个队列才恍然大悟。解决办法就是让队列名带上实例唯一标识。第二个坑手动ACK与自动ACK混用导致消息重复消费。如果你在application.yml里配了acknowledge-mode: manual但监听器方法里用了channel.basicAck同时某些分支代码里忘了确认RabbitMQ会一直重发导致客户端收到重复推送。我的建议是在监听器最外层套一个try-catch-finallyfinal里确保每条消息只ack一次never只nack重试超过3次的消息。否则一旦某条消息一直处理失败它会卡住后续所有消息形成积压。第三个坑Nginx的proxy_read_timeout太短导致连接被掐断。WebSocket是长连接默认Nginx的proxy_read_timeout是60秒也就是说60秒内后端没返回任何数据Nginx就会主动断开连接。对一个偶尔推一条告警的场景来说这个默认值完全不够。需要调到3600秒以上同时前端要做心跳保活每30秒发一个PING帧或业务心跳否则仍可能被中间网络设备掐断。第四个坑心率和Token校验的冲突。我在第6节提到心跳消息走/app/heartbeat但/app前缀的消息会先经过clientInboundChannel如果你的Token校验逻辑对每个非CONNECT命令都做一次完整JWT解析高并发心跳下CPU吃掉不少。优化办法是在ChannelInterceptor里对心跳命令做轻量校验比如只检查Redis里的session是否有效不再解析JWT签名。生产环境实测这个优化能让单实例的Tps从两千提升到近万还是很可观的。第五个坑SockJS与原生WebSocket的兼容差异。如果你的前端用的是withSockJS()有些场景比如IE浏览器会走XHR轮询兜底这时推送延迟会明显变大而且Nginx的配置要额外处理/ws/**下的多个请求路径。我对新项目一律建议原生WebSocket只有确实需要兼容老浏览器才开SockJS兜底。现在回头看看整条链路客户端通过Token完成两段式认证连接任意一个Spring Boot实例STOMP订阅关系落到本地Simple Broker发送方消息进入RabbitMQ Topic交换机广播到每一个实例的独立队列每个实例消费消息后转给本地的SimpMessagingTemplate最终推送到持锁的客户端会话期间任何一步失败都有重试和补偿机制兜底。这套架构虽然没有银弹那么玄乎但至少能让我们在业务半夜告警时少接几个消息没收到的工单。
阅读完成 · 觉得有帮助?
咨询建站