第 12 章:RabbitMQ 进阶
学习目标
- 掌握 Exchange 类型和路由机制
- 实现生产者/消费者 Java 代码
- 处理消息可靠性、幂等、死信
- 了解 RabbitMQ 集群
一、RabbitMQ 核心概念
| 概念 | 含义 |
|---|---|
| Producer | 消息生产者 |
| Consumer | 消息消费者 |
| Exchange | 交换机,接收消息并路由 |
| Queue | 队列,存储消息 |
| Binding | 交换机和队列的绑定规则 |
| Routing Key | 路由键 |
| Vhost | 虚拟主机(隔离命名空间) |
二、Exchange 类型
2.1 Direct(精确路由)
text
Routing Key = "order.created" → 绑定了 "order.created" 的 Queue 收到2.2 Topic(模式匹配)
text
Routing Key = "order.paid.cn"
Binding Pattern:
"order.*" → 匹配("*" = 一个单词)
"order.#" → 匹配("#" = 零或多个单词)2.3 Fanout(广播)
无视 Routing Key,所有绑定的 Queue 都收一份。
2.4 Headers(键值匹配)
少用,性能差。
三、安装
bash
docker run -d --name rabbitmq \
-p 5672:5672 -p 15672:15672 \
-e RABBITMQ_DEFAULT_USER=admin \
-e RABBITMQ_DEFAULT_PASS=admin123 \
rabbitmq:3-management访问 http://localhost:15672(管理界面)。
⚠️ 坑 1:
guest用户默认只能 localhost 登录,Docker 启动必须改RABBITMQ_DEFAULT_USER。
四、Spring Boot 整合
依赖
xml
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>配置
yaml
spring:
rabbitmq:
host: localhost
port: 5672
username: admin
password: admin123
virtual-host: / # ⚠️ vhost 必须已存在声明 Exchange / Queue
java
@Configuration
public class RabbitConfig {
public static final String EXCHANGE = "order.exchange";
public static final String QUEUE = "order.queue";
public static final String ROUTING_KEY = "order.created";
@Bean
public TopicExchange exchange() {
return new TopicExchange(EXCHANGE);
}
@Bean
public Queue queue() {
return QueueBuilder.durable(QUEUE).build(); // durable = 重启不丢
}
@Bean
public Binding binding() {
return BindingBuilder.bind(queue()).to(exchange()).with("order.*");
}
}⚠️ 坑 2:
Queue默认durable=true但不持久化消息;重启 queue 在,消息丢。要durable+ 消息发时MessageDeliveryMode.PERSISTENT。
五、生产者
java
@Service
public class OrderProducer {
@Resource
private RabbitTemplate rabbitTemplate;
public void send(Long orderId) {
rabbitTemplate.convertAndSend(
RabbitConfig.EXCHANGE,
"order.created", // routing key
orderId,
msg -> {
msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
msg.getMessageProperties().setMessageId(UUID.randomUUID().toString());
return msg;
}
);
}
}六、消费者
java
@Component
public class OrderConsumer {
@RabbitListener(queues = RabbitConfig.QUEUE)
public void handle(Long orderId, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag)
throws IOException {
try {
log.info("处理订单:{}", orderId);
// 业务逻辑
channel.basicAck(tag, false); // 手动 ack
} catch (Exception e) {
channel.basicNack(tag, false, false); // 不重试,直接进死信
}
}
}⚠️ 坑 3:
@RabbitListener默认是auto-ack,消息发出就认为成功 —— 关掉 auto-ack,业务成功才basicAck。 ⚠️ 坑 4:@RabbitListener方法体里抛异常,默认会无限重试,需要配retryTemplate或死信队列。
七、消息可靠性
7.1 三种保护
| 阶段 | 保证手段 |
|---|---|
| 生产发送 | publisher-confirm 异步确认 |
| 队列存储 | 队列和消息都 durable + PERSISTENT |
| 消费处理 | 手动 ack + 重试机制 |
7.2 开启发布确认
yaml
spring:
rabbitmq:
publisher-confirm-type: correlated # 异步回调
publisher-returns: truejava
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory cf) {
RabbitTemplate tpl = new RabbitTemplate(cf);
tpl.setConfirmCallback((data, ack, cause) -> {
if (!ack) log.error("消息发送失败:{}", cause);
});
return tpl;
}八、幂等消费
RabbitMQ 至少一次(重复消费不可避免),业务必须幂等。
java
// 用 Redis 判重
public boolean isProcessed(String messageId) {
Boolean first = redis.opsForValue()
.setIfAbsent("msg:" + messageId, "1", Duration.ofDays(1));
return Boolean.FALSE.equals(first); // 已存在 = 重复
}⚠️ 坑 5:幂等 token 必须全局唯一(用 messageId 或业务 ID),Redis 用
SETNX保证原子。
九、死信队列(DLX)
消息被 nack、rejected、过期时,会进死信交换机,再绑个队列专门收死信:
java
@Bean
public Queue dlq() {
return QueueBuilder.durable("order.dlq").deadLetterExchange("order.dlx").build();
}死信处理 = 人工排查 + 重新入队。
十、本章小结
| 要点 | 关键 |
|---|---|
| Exchange | Direct / Topic / Fanout / Headers |
| 可靠性 | 持久化 + 手动 ack + 发布确认 |
| 幂等 | Redis SETNX 判重,messageId 全局唯一 |
| 死信 | DLX + 死信队列,排查和重投 |
| retry | x-dead-letter-exchange + retry-interval |
| Listener | 默认无限重试,要配 retry 次数 |
动手练习
- 用 Docker 启动 RabbitMQ,声明 Topic Exchange,绑定
order.*队列 - 生产者发 "order.created"、"order.paid.cn",消费者只收
order.* - 故意抛异常,观察死信队列收到消息
下一章:第 13 章:Kafka 进阶 →