RabbitMQ 集成:spring-boot-starter-amqp 收发与确认

订单创建后要发通知、异步扣积分、同步给下游——消息队列解耦了这些链路。RabbitMQ 核心模型:生产者把消息发到交换机(Exchange),交换机按 Binding 投递到队列(Queue),消费者从队列取消息。Spring Boot 用 spring-boot-starter-amqp 封装收发。

依赖与连接配置

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
spring:
  rabbitmq:
    host: 127.0.0.1
    port: 5672
    username: guest
    password: guest
    publisher-confirm-type: correlated   # 生产者确认:消息到达交换机
    publisher-returns: true              # 路由不到队列时回退给生产者
    listener:
      simple:
        acknowledge-mode: manual         # 消费者手动确认(见下)

声明队列、交换机、绑定

@Configuration
public class RabbitConfig {
    public static final String EXCHANGE = "order.ex";     // direct 直连交换机
    public static final String QUEUE   = "order.created.q";
    public static final String ROUTING = "order.created";

    @Bean
    Queue queue() { return QueueBuilder.durable(QUEUE).build(); }   // durable:重启不丢

    @Bean
    DirectExchange exchange() { return new DirectExchange(EXCHANGE); }

    @Bean
    Binding binding() {
        return BindingBuilder.bind(queue()).to(exchange()).with(ROUTING);
    }
}

交换机类型:direct(精确 routingKey)、topic(通配如 order.*)、fanout(广播),按场景选用。

发送:RabbitTemplate

@Service
public class OrderPublisher {
    private final RabbitTemplate rabbitTemplate;   // 构造注入

    public void publish(OrderCreatedEvent event) {
        rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE, RabbitConfig.ROUTING, event);
        // 输出:消息 → 交换机 → 按 binding 进入 order.created.q
    }
}

声明 JSON 转换器后对象自动序列化为 JSON(收发两端必须一致):

@Bean
public MessageConverter messageConverter() {
    return new Jackson2JsonMessageConverter();   // 对象→JSON,消费端 JSON→对象
}

消费与手动确认

@Component
public class OrderCreatedListener {

    @RabbitListener(queues = RabbitConfig.QUEUE)
    public void onMessage(OrderCreatedEvent event, Channel channel,
                          @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        try {
            handle(event);                       // 业务处理(幂等见第 41 章)
            channel.basicAck(tag, false);        // 确认成功,Broker 删除消息
        } catch (Exception e) {
            channel.basicNack(tag, false, true); // 拒绝并重回队列重试
            // 多次失败建议 requeue=false + 死信队列,避免消息原地打转
        }
    }
}

手动确认(acknowledge-mode: manual)把"删消息"交给业务:成功才 ack,失败可 nack 重投。若不需精确控制,可改回默认 auto 并配 spring.rabbitmq.listener.simple.retry.enabled=true 让框架自动重试(次数耗尽后默认拒绝且不重回队列)。

可靠性三板斧

  1. 生产者确认:publisher-confirm-type=correlated,配合 ConfirmCallback/ReturnsCallback 感知消息是否到达交换机、是否可路由,失败则记录重发;
  2. 持久化:durable 队列 + 消息持久化,Broker 重启不丢;
  3. 消费者手动确认:处理成功才 ack,防止"消费即删、业务却没做成"。

小结:starter-amqp 下 RabbitTemplate 发、@RabbitListener 收;队列/交换机声明式配置;可靠性靠"生产者确认 + 持久化 + 手动 ack + 死信兜底"闭环。

笔记加载中…