Spring Boot STOMP消费者无法接收RabbitMQ中消息的问题排查与解决
Spring Boot STOMP消费者无法接收RabbitMQ中消息的问题排查与解决
为什么你的消费者收不到消息?
我先帮你梳理下核心问题所在:
@MessageMapping的用途误解:你在消费者中使用的@MessageMapping("/topic/public.messages")是用来处理前端WebSocket客户端发送到该目的地的消息,而非订阅RabbitMQ Broker上的Topic消息。这个注解属于Spring WebSocket的客户端→服务端消息处理机制,不适用服务端→服务端的Broker消息订阅场景。- 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); // 处理消息的业务逻辑 } }
验证效果
配置完成后重启消费者服务:
- 打开RabbitMQ UI,你会看到对应队列(方案2的
public.messages.queue或方案1的临时队列)已绑定到amq.topic交换器。 - 启动生产者服务后,消费者控制台会开始打印接收到的消息日志。
内容来源于stack exchange
相关产品推荐
相关产品推荐

