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

Spring Boot STOMP消费者无法接收RabbitMQ中消息的问题排查与解决

Spring Boot STOMP消费者无法接收RabbitMQ中消息的问题排查与解决

为什么你的消费者收不到消息?

我先帮你梳理下核心问题所在:

  1. @MessageMapping的用途误解:你在消费者中使用的@MessageMapping("/topic/public.messages")是用来处理前端WebSocket客户端发送到该目的地的消息,而非订阅RabbitMQ Broker上的Topic消息。这个注解属于Spring WebSocket的客户端→服务端消息处理机制,不适用服务端→服务端的Broker消息订阅场景。
  2. RabbitMQ中无队列绑定:你在RabbitMQ UI中看到amq.topic交换器有消息流入,但没有任何队列绑定到这个交换器。RabbitMQ的规则是:如果交换器收到的消息没有匹配的队列,会直接被丢弃(除非配置了死信交换器)。所以生产者发送的消息其实都被RabbitMQ丢弃了,消费者自然收不到任何内容。

可行解决方案

根据你的需求(Spring Boot微服务作为STOMP消息的服务端消费者),推荐两种实践方案:

方案1:使用RabbitMQ STOMP客户端直接订阅Topic

这种方式让消费者服务通过STOMP协议直接连接RabbitMQ,自动创建临时队列并绑定到amq.topic,从而接收消息。

步骤1:添加依赖

在消费者的pom.xml中引入RabbitMQ STOMP客户端依赖:

<dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>rabbitmq-stomp-client</artifactId>
</dependency>
步骤2:实现STOMP消费者组件

创建Spring组件,初始化STOMP连接并订阅目标Topic:

import com.rabbitmq.stomp.client.StompClient;
import com.rabbitmq.stomp.client.StompConnection;
import com.rabbitmq.stomp.client.StompFrame;
import com.rabbitmq.stomp.client.StompFrameHandler;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import org.springframework.stereotype.Component;

@Component
public class StompTopicConsumer {

    private StompConnection stompConnection;

    @PostConstruct
    public void initStompConnection() throws Exception {
        // 初始化STOMP客户端,连接RabbitMQ的STOMP端口(61613)
        StompClient stompClient = new StompClient();
        stompClient.setHost("localhost");
        stompClient.setPort(61613);
        stompClient.setLogin("guest");
        stompClient.setPasscode("guest");

        // 建立连接
        stompConnection = stompClient.connect();
        System.out.println("Connected to RabbitMQ STOMP broker");

        // 订阅目标Topic
        stompConnection.subscribe("/topic/public.messages", new StompFrameHandler() {
            @Override
            public void handleFrame(StompFrame frame) {
                String messageContent = new String(frame.getBody());
                System.out.println("Received STOMP message: " + messageContent);
                // 在这里添加你的消息处理逻辑
            }
        });
        System.out.println("Subscribed to /topic/public.messages");
    }

    @PreDestroy
    public void cleanupConnection() throws Exception {
        if (stompConnection != null && stompConnection.isConnected()) {
            stompConnection.disconnect();
            System.out.println("Disconnected from STOMP broker");
        }
    }
}

方案2:使用RabbitMQ AMQP监听队列(更稳定的服务端消费方式)

如果你的消费者服务不需要严格依赖STOMP协议,直接通过AMQP协议监听RabbitMQ队列是Spring Boot微服务的常用实践,稳定性和可维护性更强。

步骤1:添加RabbitMQ AMQP依赖
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
步骤2:配置队列与交换器绑定

创建RabbitMQ配置类,将自定义队列绑定到RabbitMQ自带的amq.topic交换器:

import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.TopicExchange;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitMqConfig {

    // 定义持久化队列
    @Bean
    public Queue publicMessagesQueue() {
        return new Queue("public.messages.queue", true);
    }

    // 引用RabbitMQ自带的amq.topic交换器
    @Bean
    public TopicExchange amqTopicExchange() {
        return new TopicExchange("amq.topic");
    }

    // 将队列绑定到amq.topic,路由键匹配public.messages(对应STOMP的/topic/public.messages)
    @Bean
    public Binding bindingPublicMessages(Queue publicMessagesQueue, TopicExchange amqTopicExchange) {
        return BindingBuilder.bind(publicMessagesQueue).to(amqTopicExchange).with("public.messages");
    }
}
步骤3:创建消息监听器

用@RabbitListener注解监听队列,接收并处理消息:

import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
public class AmqpMessageListener {

    @RabbitListener(queues = "public.messages.queue")
    public void handlePublicMessage(String message) {
        System.out.println("Received AMQP message from queue: " + message);
        // 处理消息的业务逻辑
    }
}

验证效果

配置完成后重启消费者服务:

  1. 打开RabbitMQ UI,你会看到对应队列(方案2的public.messages.queue或方案1的临时队列)已绑定到amq.topic交换器。
  2. 启动生产者服务后,消费者控制台会开始打印接收到的消息日志。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 09:59:52