Spring Boot基于RabbitMQ实现WebSocket会话专属动态订阅的方案咨询
基于RabbitMQ的动态WebSocket会话消息路由实现方案
需求概述
- 每个WebSocket会话订阅专属通道,仅接收对应通道的消息
- 握手时通过
conversationId参数指定会话关联的对话ID,仅接收该ID的消息 - REST接口接收消息后,路由到对应WebSocket会话
- 多实例部署时,确保消息仅被持有目标会话的实例接收
初始尝试方案
1. 握手阶段保存会话属性
@Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception { String conversationId = ((ServletServerHttpRequest) request).getServletRequest().getParameter("conversationId"); attributes.put("conversationId", conversationId); return true; }
2. 会话生命周期管理
@Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { String conversationId = session.getAttributes().get("conversationId").toString(); sessionStore.sessions().put(session.getId(), session); Queue q = new Queue(session.getId()); admin.declareQueue(q); container.addQueues(q); sessionStore.queues().put(session.getId(), q); Binding b = BindingBuilder.bind(q).to(exchange).with("conversation." + conversationId); admin.declareBinding(b); sessionStore.bindings().put(session.getId(), b); } @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { Binding b = sessionStore.bindings().remove(session.getId()); admin.removeBinding(b); Queue q = sessionStore.queues().remove(session.getId()); container.removeQueues(q); admin.deleteQueue(q.getName()); sessionStore.sessions().remove(session.getId()); }
3. 单容器消息监听
@Bean public DirectMessageListenerContainer container(ConnectionFactory connectionFactory, WebSocketSessionStore sessionStore) { DirectMessageListenerContainer container = new DirectMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setMessageListener(new MessageListener() { @Override public void onMessage(Message message) { WebSocketSession ws = sessionStore.sessions().get(message.getMessageProperties().getConsumerQueue()); try { ws.sendMessage(new TextMessage(message.getBody())); } catch (IOException e) { throw new RuntimeException(e); } } }); return container; }
当前遇到的问题
- 单容器需通过HashMap匹配队列与会话,管理复杂度高
- 不确定为每个会话创建独立容器/监听器是否为最佳实践,以及资源开销情况
- 不明确
RabbitListenerEndpoint的作用与适用场景
优化后的实现方案
核心思路
利用Spring AMQP的RabbitListenerEndpointRegistry动态注册会话专属消费者,每个会话对应独立的监听逻辑,同时通过RabbitMQ的自动删除队列特性简化资源清理。
1. 配置基础RabbitMQ组件
@Bean public DirectExchange conversationExchange() { // 持久化Exchange,确保重启后不丢失 return new DirectExchange("conversation.exchange", true, false); } @Bean public SimpleRabbitListenerContainerFactory listenerContainerFactory(ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); factory.setAutoStartup(false); // 手动控制启动时机,避免提前消费 return factory; }
2. WebSocket会话管理与动态消费者注册
@Component public class ConversationWebSocketHandler extends TextWebSocketHandler { private final AmqpAdmin amqpAdmin; private final DirectExchange conversationExchange; private final RabbitListenerEndpointRegistry endpointRegistry; private final SimpleRabbitListenerContainerFactory containerFactory; // 会话与端点ID映射,用于关闭时清理 private final Map<String, String> sessionEndpointMap = new ConcurrentHashMap<>(); // 会话缓存,用于消息推送 private final Map<String, WebSocketSession> sessionMap = new ConcurrentHashMap<>(); public ConversationWebSocketHandler(AmqpAdmin amqpAdmin, DirectExchange conversationExchange, RabbitListenerEndpointRegistry endpointRegistry, SimpleRabbitListenerContainerFactory containerFactory) { this.amqpAdmin = amqpAdmin; this.conversationExchange = conversationExchange; this.endpointRegistry = endpointRegistry; this.containerFactory = containerFactory; } @Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception { String conversationId = ((ServletServerHttpRequest) request).getServletRequest().getParameter("conversationId"); if (StringUtils.isBlank(conversationId)) { response.setStatusCode(HttpStatus.BAD_REQUEST); return false; } attributes.put("conversationId", conversationId); return true; } @Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { String sessionId = session.getId(); String conversationId = session.getAttributes().get("conversationId").toString(); sessionMap.put(sessionId, session); // 创建会话专属队列:自动删除、非持久化,会话关闭后RabbitMQ自动清理 Queue sessionQueue = new Queue(sessionId, false, false, true); amqpAdmin.declareQueue(sessionQueue); // 绑定队列到Exchange,路由键为conversation.{conversationId} Binding binding = BindingBuilder.bind(sessionQueue) .to(conversationExchange) .with("conversation." + conversationId); amqpAdmin.declareBinding(binding); // 动态创建RabbitListener端点 SimpleRabbitListenerEndpoint endpoint = new SimpleRabbitListenerEndpoint(); String endpointId = sessionId + "-conversation-listener"; endpoint.setId(endpointId); endpoint.setQueueNames(sessionId); // 专属消息监听逻辑,直接关联当前会话 endpoint.setMessageListener(message -> { try { WebSocketSession wsSession = sessionMap.get(sessionId); if (wsSession != null && wsSession.isOpen()) { wsSession.sendMessage(new TextMessage(message.getBody())); } } catch (IOException e) { // 异常处理:记录日志、主动关闭会话等 throw new RuntimeException("推送消息到WebSocket会话失败", e); } }); // 注册端点并启动容器 MessageListenerContainer container = endpointRegistry.registerListenerContainer(endpoint, containerFactory); container.start(); sessionEndpointMap.put(sessionId, endpointId); } @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { String sessionId = session.getId(); sessionMap.remove(sessionId); // 停止并注销RabbitMQ消费者 String endpointId = sessionEndpointMap.remove(sessionId); if (endpointId != null) { MessageListenerContainer container = endpointRegistry.getListenerContainer(endpointId); if (container != null) { container.stop(); endpointRegistry.unregisterListenerContainer(endpointId); } } // 手动删除队列(可选,因为队列已设置自动删除) amqpAdmin.deleteQueue(sessionId); } }
3. REST接口消息路由实现
@RestController @RequestMapping("/conversations") public class ConversationMessageController { private final RabbitTemplate rabbitTemplate; private final DirectExchange conversationExchange; public ConversationMessageController(RabbitTemplate rabbitTemplate, DirectExchange conversationExchange) { this.rabbitTemplate = rabbitTemplate; this.conversationExchange = conversationExchange; } @PostMapping("/{conversationId}") public ResponseEntity<Void> sendMessage(@PathVariable String conversationId, @RequestBody String message) { // 发送消息到指定路由键,RabbitMQ会路由到所有绑定该键的队列 rabbitTemplate.convertAndSend(conversationExchange.getName(), "conversation." + conversationId, message); return ResponseEntity.ok().build(); } }
方案优势
- 逻辑清晰:每个会话对应独立的消费者逻辑,无需全局HashMap匹配队列与会话
- 资源高效:Spring AMQP的容器工厂会复用连接与线程池,单个会话的资源开销极低
- 多实例兼容:会话专属队列仅在创建它的实例上有消费者,消息只会被持有目标会话的实例接收
- 自动清理:队列设置为自动删除,即使服务异常退出,RabbitMQ也会自动清理无消费者的队列
问题解答
- 单容器管理复杂度:通过动态注册
RabbitListenerEndpoint,每个会话的消费逻辑独立,彻底避免全局映射的管理成本 - 单会话容器开销:Spring AMQP的容器工厂会复用底层连接与线程,实际资源消耗远低于手动创建多个独立容器
- RabbitListenerEndpoint的作用:是Spring AMQP提供的动态消费者注册抽象,替代静态
@RabbitListener注解,完美适配运行时动态创建队列的场景
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

