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

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();
        }
    }
}

额外排查步骤

  1. 登录RabbitMQ管理后台(默认端口15672),检查:
    • 交换机restaurant-score-exchange是否存在
    • 队列restaurant-score-queue是否存在,且与交换机正确绑定
    • 队列的Ready消息数是否有增长,若增长但未消费,说明消费者监听队列错误或连接异常
  2. 查看消费者微服务日志,排查是否有RabbitMQ连接失败、队列不存在等异常
  3. 确认两个微服务连接的是同一个RabbitMQ实例,host、port、账号密码配置一致

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 19:55:33