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 让框架自动重试(次数耗尽后默认拒绝且不重回队列)。
可靠性三板斧
- 生产者确认:publisher-confirm-type=correlated,配合 ConfirmCallback/ReturnsCallback 感知消息是否到达交换机、是否可路由,失败则记录重发;
- 持久化:durable 队列 + 消息持久化,Broker 重启不丢;
- 消费者手动确认:处理成功才 ack,防止"消费即删、业务却没做成"。
小结:starter-amqp 下 RabbitTemplate 发、@RabbitListener 收;队列/交换机声明式配置;可靠性靠"生产者确认 + 持久化 + 手动 ack + 死信兜底"闭环。