RabbitMQ消息收发异常排查:消息已发送但未被处理
RabbitMQ消费者未触发的问题排查与修复
核心问题分析
- 消费者监听对象错误:
@RabbitListener(queues = "restaurant-score-exchange")里填的是交换机名称,但RabbitMQ消费者只能监听队列,交换机仅负责路由消息到队列,无法直接被消费者监听。 - 缺少队列与交换机绑定配置:你的配置类只定义了交换机,没有创建队列,也未将队列与交换机通过路由键绑定,导致生产者发送的消息没有目标队列可存储,自然无法被消费者接收。
- 消息格式不匹配:生产者发送纯字符串
"newUser",但消费者尝试将其解析为JSON对象,即便消息能送达,也会触发解析异常,中断逻辑执行。 - 控制器注入写法错误:
@Autowired private final RabbitMQMessageProducer rabbitMQMessageProducer;中,final字段搭配@Autowired会导致注入失败,推荐用构造器注入。 - 冗余的
@Import:若Spring Boot已通过组件扫描加载RabbitMQConfig,生产者和消费者类上的@Import(RabbitMQConfig.class)属于重复配置,可能引发Bean冲突。
修正后的代码实现
1. 完善RabbitMQ配置类(创建队列+绑定交换机)
@Configuration public class RabbitMQConfig { // 统一定义常量,避免硬编码 public static final String RESTAURANT_SCORE_EXCHANGE = "restaurant-score-exchange"; public static final String RESTAURANT_SCORE_QUEUE = "restaurant-score-queue"; @Bean public DirectExchange restaurantScoreExchange() { return new DirectExchange(RESTAURANT_SCORE_EXCHANGE); } @Bean public Queue restaurantScoreQueue() { // 创建持久化队列,避免服务重启丢失消息 return QueueBuilder.durable(RESTAURANT_SCORE_QUEUE).build(); } @Bean public Binding binding(Queue queue, DirectExchange exchange) { // 将队列绑定到交换机,指定路由键规则 return BindingBuilder.bind(queue).to(exchange).with("restaurant.score.update"); } }
2. 修正生产者代码
@Component public class RabbitMQMessageProducer { private final AmqpTemplate rabbitTemplate; private final ObjectMapper objectMapper; // 构造器注入(Spring 4.3+无需@Autowired) public RabbitMQMessageProducer(AmqpTemplate rabbitTemplate, ObjectMapper objectMapper) { this.rabbitTemplate = rabbitTemplate; this.objectMapper = objectMapper; } public void sendUpdateMessage(Long restaurantId, String message) { // 构造符合消费者解析逻辑的JSON消息 ObjectNode messageNode = objectMapper.createObjectNode(); messageNode.put("restaurantId", restaurantId); messageNode.put("eventType", message); rabbitTemplate.convertAndSend( RabbitMQConfig.RESTAURANT_SCORE_EXCHANGE, "restaurant.score.update", objectMapper.writeValueAsString(messageNode) ); } }
3. 修正控制器注入
// 推荐构造器注入方式 private final RabbitMQMessageProducer rabbitMQMessageProducer; public YourController(RabbitMQMessageProducer rabbitMQMessageProducer) { this.rabbitMQMessageProducer = rabbitMQMessageProducer; } void randomMethod() { ... rabbitMQMessageProducer.sendUpdateMessage(newUser.getFavoriteRestaurantId(), "newUser"); ... }
4. 修正消费者代码
@Component public class RabbitMQMessageReceiver { private final ObjectMapper objectMapper; public RabbitMQMessageReceiver(ObjectMapper objectMapper) { this.objectMapper = objectMapper; } // 监听正确的队列名称 @RabbitListener(queues = RabbitMQConfig.RESTAURANT_SCORE_QUEUE) public void handleMessage(String message) { try { JsonNode messageJson = objectMapper.readTree(message); Long restaurantId = messageJson.get("restaurantId").asLong(); String eventType = messageJson.get("eventType").asText(); // 执行餐厅评分更新逻辑 System.out.println("处理餐厅[" + restaurantId + "]的评分更新,事件类型:" + eventType); } catch (Exception e) { System.err.println("消息处理失败:" + e.getMessage()); e.printStackTrace(); } } }
额外排查步骤
- 登录RabbitMQ管理后台(默认端口15672),检查:
- 交换机
restaurant-score-exchange是否存在 - 队列
restaurant-score-queue是否存在,且与交换机正确绑定 - 队列的
Ready消息数是否有增长,若增长但未消费,说明消费者监听队列错误或连接异常
- 交换机
- 查看消费者微服务日志,排查是否有RabbitMQ连接失败、队列不存在等异常
- 确认两个微服务连接的是同一个RabbitMQ实例,host、port、账号密码配置一致
内容的提问来源于stack exchange,提问作者spozzi
相关产品推荐
相关产品推荐

