You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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也会自动清理无消费者的队列

问题解答

  1. 单容器管理复杂度:通过动态注册RabbitListenerEndpoint,每个会话的消费逻辑独立,彻底避免全局映射的管理成本
  2. 单会话容器开销:Spring AMQP的容器工厂会复用底层连接与线程,实际资源消耗远低于手动创建多个独立容器
  3. RabbitListenerEndpoint的作用:是Spring AMQP提供的动态消费者注册抽象,替代静态@RabbitListener注解,完美适配运行时动态创建队列的场景

内容的提问来源于stack exchange,提问作者Paul

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 15:42:05