已发布消息但RabbitMQ队列仍为空的问题求助
RabbitMQ队列显示为空但消费者已收到消息的问题排查
我正在开发基于RabbitMQ的Spring应用,按教程配置后出现异常:通过Postman或RabbitMQ管理端(localhost:15672)发布消息时,系统提示消息已发布,应用控制台也显示消费者成功收到消息,但RabbitMQ管理端的队列页面始终显示队列为空。以下是完整的代码配置,请帮忙排查原因:
配置类(RabbitMQConfig)
package ro.tuc.ds2020.config; import org.springframework.amqp.core.*; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitMQConfig { @Value("${rabbitmq.queue.name}") private String queue; @Value("${rabbitmq.queue.json.name}") private String jsonQueue; @Value("${rabbitmq.queue.exchange}") private String exchange; @Value("${rabbitmq.queue.routing_key_one}") private String routingKeyOne; @Value("${rabbitmq.queue.routing_key_json}") private String routingKeyJson; @Bean public Queue queue() { return new Queue(queue); } @Bean public Queue jsonQueue() { return new Queue(jsonQueue, true); } @Bean public TopicExchange exchange() { return new TopicExchange(exchange, true, false); } @Bean public Binding binding() { return BindingBuilder.bind(queue()).to(exchange()).with(routingKeyOne); } @Bean public Binding jsonBinding() { return BindingBuilder.bind(jsonQueue()).to(exchange()).with(routingKeyJson); } @Bean public MessageConverter converter() { return new Jackson2JsonMessageConverter(); } @Bean public AmqpTemplate amqpTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); rabbitTemplate.setMessageConverter(converter()); return rabbitTemplate; } }
消费者类(RabbitMQJsonConsumer)
package ro.tuc.ds2020.consumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Service; import ro.tuc.ds2020.dtos.MeasurementDTO; @Service public class RabbitMQJsonConsumer { private static final Logger LOGGER = LoggerFactory.getLogger(RabbitMQJsonConsumer.class); @RabbitListener(queues = {"${rabbitmq.queue.json.name}"}) public void consumeJsonMessage(MeasurementDTO measurementDTO) { LOGGER.info(String.format("Received JSON message here -> %s", measurementDTO.toString())); } }
控制器类(MessageJsonController)
package ro.tuc.ds2020.controllers; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.*; import ro.tuc.ds2020.dtos.MeasurementDTO; import ro.tuc.ds2020.publisher.RabbitMQJsonProducer; @RequestMapping("/api/v1") @RestController @CrossOrigin(origins = "http://localhost:4200", allowCredentials = "true") public class MessageJsonController { private RabbitMQJsonProducer jsonProducer; public MessageJsonController(RabbitMQJsonProducer rabbitMQJsonProducer) { this.jsonProducer = rabbitMQJsonProducer; } @PostMapping("/publish") public ResponseEntity<String> sendJsonMessage(@RequestBody MeasurementDTO measurementDTO) { jsonProducer.sendJsonMessage(measurementDTO); return ResponseEntity.ok("Json message sent to RabbitMQ ..."); } }
生产者类(RabbitMQJsonProducer)
package ro.tuc.ds2020.publisher; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import ro.tuc.ds2020.dtos.MeasurementDTO; @Service public class RabbitMQJsonProducer { @Value("${rabbitmq.queue.exchange}") private String exchange; @Value("${rabbitmq.queue.routing_key_json}") private String routingKeyJson; private static final Logger LOGGER = LoggerFactory.getLogger(RabbitMQJsonProducer.class); private RabbitTemplate rabbitTemplate; @Autowired public RabbitMQJsonProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } public void sendJsonMessage(MeasurementDTO measurementDTO) { LOGGER.info(String.format("Json message sent -> %s", measurementDTO.toString())); rabbitTemplate.convertAndSend(exchange, routingKeyJson, measurementDTO); } }
配置文件(application.properties)
spring.rabbitmq.host = localhost spring.rabbitmq.port = 5672 spring.rabbitmq.username = guest spring.rabbitmq.password = guest rabbitmq.queue.name = queue_1 rabbitmq.queue.json.name = queue_json rabbitmq.queue.exchange = exchange rabbitmq.queue.routing_key_one = routing_key_1 rabbitmq.queue.routing_key_json = routing_key_json
问题原因分析
这种现象是消息被消费者立即消费并确认导致的,属于正常的RabbitMQ工作流程:
- Spring AMQP的
@RabbitListener默认采用AUTO确认模式,消息到达队列后,消费者会立即获取消息并自动向RabbitMQ发送确认信号。 - 由于你的消费逻辑处理速度极快,消息刚进入队列就被取走并确认,RabbitMQ管理端来不及捕获到队列中有消息的状态,因此显示为空。
- 控制台能收到消费日志,说明交换机、路由键、队列的绑定关系完全正常,消息没有丢失或进入死信队列。
验证与解决方案
验证方法
- 临时停止消费者服务,再通过Postman或管理端发送消息,此时查看RabbitMQ管理端的队列,能看到消息堆积。
- 重启消费者服务,消息会被立即消费,队列再次变为空,即可确认是消费速度快导致的显示问题。
调整消费确认模式(可选)
如果需要让消息在队列中停留一段时间,或手动控制确认时机,可以修改@RabbitListener的确认模式为手动确认:
import com.rabbitmq.client.Channel; import org.springframework.amqp.core.AmqpHeaders; import org.springframework.messaging.handler.annotation.Header; // ... @RabbitListener(queues = {"${rabbitmq.queue.json.name}"}, ackMode = "MANUAL") public void consumeJsonMessage(MeasurementDTO measurementDTO, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException { LOGGER.info(String.format("Received JSON message here -> %s", measurementDTO.toString())); // 处理完业务逻辑后手动确认消息 channel.basicAck(tag, false); // 若处理失败,可拒绝消息并重新入队 // channel.basicNack(tag, false, true); }
修改后,消息会留在队列中直到你调用basicAck手动确认,此时管理端就能看到消息存在的状态。
队列持久化检查
你的jsonQueue已配置持久化(new Queue(jsonQueue, true)),即使重启RabbitMQ,未消费的消息也不会丢失,这部分配置是正确的。
内容的提问来源于stack exchange,提问作者Cornel
相关产品推荐
相关产品推荐

