第 62 章:消息队列与 RabbitMQ
学习目标
- 理解 MQ 的核心价值:异步、解耦、削峰
- 掌握 RabbitMQ 五大模式与 Spring 集成
- 学会可靠投递、幂等消费、死信队列
一、为什么需要 MQ?
MQ 的三大核心价值:
| 价值 | 含义 |
|---|---|
| 异步 | 写完就走,下游慢慢处理 |
| 解耦 | 服务之间通过消息通信,不依赖接口契约 |
| 削峰 | 流量高峰先堆积在 MQ,系统按能力消费 |
二、MQ 选型对比
| 维度 | RabbitMQ | RocketMQ | Kafka |
|---|---|---|---|
| 吞吐量 | 万级 | 十万级 | 百万级 |
| 延迟 | 微秒级 | 毫秒级 | 毫秒级 |
| 消息可靠性 | 强 | 强 | 中(默认可能丢) |
| 事务消息 | 弱(要事务机制) | 强(原生支持) | 弱 |
| 生态 | 成熟 | 阿里系 | 大数据、流计算 |
| 适用 | 企业业务、可靠投递 | 电商、金融(阿里主推) | 日志、流计算 |
三、RabbitMQ 核心概念
| 概念 | 说明 |
|---|---|
| Producer | 消息生产者 |
| Consumer | 消息消费者 |
| Exchange | 交换机,接收生产者的消息并路由 |
| Queue | 队列,真正存储消息的地方 |
| Binding | 交换机和队列之间的绑定规则 |
| Routing Key | 路由键,交换机用它决定消息去哪个队列 |
| VHost | 虚拟主机,多租户隔离 |
四、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: ${RABBITMQ_PASSWORD}
virtual-host: /taskflow
publisher-confirm-type: correlated # ① 发送确认
publisher-returns: true # ② 路由失败回调
listener:
simple:
acknowledge-mode: manual # ③ 手动 ack
prefetch: 10 # ④ 每次拉取 10 条
concurrency: 4 # ⑤ 4 个消费者线程
max-concurrency: 8
retry:
enabled: true
max-attempts: 3
initial-interval: 1s五、声明交换机和队列
java
@Configuration
public class RabbitConfig {
// ① 业务交换机(Topic 模式)
@Bean
public TopicExchange orderExchange() {
return ExchangeBuilder.topicExchange("order.exchange")
.durable(true)
.build();
}
// ② 下单队列
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
.deadLetterExchange("dlx.exchange") // 死信交换机
.deadLetterRoutingKey("order.dlq")
.ttl(24 * 60 * 60 * 1000) // 24 小时过期
.maxLength(100000) // 最大消息数
.build();
}
// ③ 绑定:order.* 路由到 order.queue
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(orderExchange())
.with("order.#");
}
// ④ 死信交换机
@Bean
public DirectExchange dlxExchange() {
return ExchangeBuilder.directExchange("dlx.exchange").build();
}
@Bean
public Queue dlq() {
return QueueBuilder.durable("order.dlq").build();
}
@Bean
public Binding dlqBinding() {
return BindingBuilder.bind(dlq()).to(dlxExchange()).with("order.dlq");
}
}六、生产者:可靠发送
java
@Component
@RequiredArgsConstructor
@Slf4j
public class OrderProducer {
private final RabbitTemplate rabbitTemplate;
// ① 发送消息(带 ConfirmCallback)
public void send(OrderMessage message) {
// ② 业务唯一 ID(用于幂等)
String msgId = UUID.randomUUID().toString();
// ③ 构造消息
Message msg = MessageBuilder
.withBody(JSON.toJSONBytes(message))
.setContentType(MessageProperties.CONTENT_TYPE_JSON)
.setMessageId(msgId) // 幂等 key
.setCorrelationId(message.getOrderId().toString())
.build();
// ④ 设置 ConfirmCallback(发送成功回调)
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
if (ack) {
log.info("消息发送成功 id={}", msgId);
} else {
log.error("消息发送失败 id={} cause={}", msgId, cause);
// 重试或入库
}
});
// ⑤ 设置 ReturnsCallback(路由失败回调)
rabbitTemplate.setReturnsCallback(returned ->
log.error("消息路由失败 exchange={} key={}",
returned.getExchange(), returned.getRoutingKey()));
// ⑥ 发送
CorrelationData corr = new CorrelationData(msgId);
rabbitTemplate.convertAndSend("order.exchange", "order.created", msg, corr);
}
}本地消息表 + MQ(生产级可靠投递)
java
@Transactional
public void placeOrder(OrderDTO dto) {
// ① 业务落库
orderMapper.insert(order);
// ② 同一事务里插入本地消息表
LocalMessage msg = new LocalMessage();
msg.setId(UUID.randomUUID().toString());
msg.setTopic("order.created");
msg.setPayload(JSON.toJSONString(order));
msg.setStatus(0); // 待发送
localMessageMapper.insert(msg);
}
// ③ 后台任务:扫表发送
@Scheduled(fixedDelay = 1000)
public void sendPendingMessages() {
List<LocalMessage> msgs = localMessageMapper.selectPending(100);
for (LocalMessage msg : msgs) {
try {
rabbitTemplate.convertAndSend("order.exchange",
"order.created", msg.getPayload());
msg.setStatus(1); // 已发送
localMessageMapper.updateById(msg);
} catch (Exception e) {
log.warn("消息发送失败,稍后重试 id={}", msg.getId());
}
}
}七、消费者:手动 ack + 幂等
java
@Component
@RequiredArgsConstructor
@Slf4j
@RocketMQMessageListener // 伪注解示意,实际用 @RabbitListener
public class OrderConsumer {
private final InventoryService inventoryService;
private final RedisTemplate<String, Object> redis;
// ① 监听队列
@RabbitListener(queues = "order.queue")
public void onMessage(Message message, Channel channel) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
String msgId = message.getMessageProperties().getMessageId();
try {
// ② 幂等检查(关键!)
String key = "mq:consumed:" + msgId;
Boolean firstTime = redis.opsForValue().setIfAbsent(key, "1", 24, TimeUnit.HOURS);
if (Boolean.FALSE.equals(firstTime)) {
log.warn("重复消息,直接 ack msgId={}", msgId);
channel.basicAck(deliveryTag, false); // 重复消息直接 ack
return;
}
// ③ 业务处理
OrderMessage order = JSON.parseObject(message.getBody(), OrderMessage.class);
inventoryService.deduct(order.getSkuId(), order.getQuantity());
// ④ 手动 ack
channel.basicAck(deliveryTag, false);
log.info("消息处理成功 msgId={}", msgId);
} catch (BusinessException e) {
// ⑤ 业务异常:消息重试(可入死信)
log.warn("业务异常,重试 msgId={}", msgId, e);
channel.basicNack(deliveryTag, false, true);
} catch (Exception e) {
// ⑥ 系统异常:消息重试或入死信
log.error("系统异常 msgId={}", msgId, e);
channel.basicNack(deliveryTag, false, false); // 不重试 → 进入死信
}
}
}幂等的几种实现
java
// 方案 1:Redis SETNX(最常用)
Boolean ok = redis.opsForValue().setIfAbsent(
"mq:msg:" + msgId, "1", 24, TimeUnit.HOURS);
if (!ok) return; // 已消费过
// 方案 2:数据库唯一索引
// 消息处理表加唯一约束 msg_id,重复消息 INSERT 失败即视为已处理
INSERT INTO mq_consumed_log (msg_id, create_time) VALUES (#{msgId}, NOW());
// DuplicateKeyException → 已处理
// 方案 3:业务判断
// 订单已处理过 status=PAID,跳过
if (order.getStatus() == PAID) return;八、五大消息模式
1. 简单模式(Hello World)
java
// 生产者
rabbitTemplate.convertAndSend("hello.queue", "Hello World");
// 消费者
@RabbitListener(queues = "hello.queue")
public void onMessage(String msg) {
System.out.println("收到:" + msg);
}2. Work Queue(任务队列)
多个消费者分摊任务(轮询 / 公平分发)
3. Publish/Subscribe(发布订阅)
java
// Fanout 交换机:忽略 routing key,广播到所有绑定队列
@Bean
public FanoutExchange fanout() { return new FanoutExchange("notify.exchange"); }4. Routing(路由)
java
// Direct 交换机:完全匹配 routing key
rabbitTemplate.convertAndSend("log.exchange", "error", "数据库连接失败");5. Topic(主题,最常用)
java
// Topic 交换机:通配符匹配
// * 匹配一个单词,# 匹配多个单词
rabbitTemplate.convertAndSend("order.exchange", "order.vip.created", msg); // 匹配 # 和 *
rabbitTemplate.convertAndSend("order.exchange", "order.normal.paid", msg); // 只匹配 #九、死信队列(DLX)
死信:消息被拒绝、过期、队列满时的归宿。
java
// 声明队列时指定 DLX
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
.deadLetterExchange("dlx.exchange")
.deadLetterRoutingKey("order.dlq")
.build();
}java
@Component
@Slf4j
public class DlqConsumer {
@RabbitListener(queues = "order.dlq")
public void onDead(Message msg) {
log.error("死信消息: msgId={} body={}",
msg.getMessageProperties().getMessageId(),
new String(msg.getBody()));
// 告警 / 入库 / 人工处理
}
}十、消息可靠性保障总结
| 阶段 | 风险 | 保障 |
|---|---|---|
| 生产者发送 | 网络抖动、MQ 宕机 | ConfirmCallback + 本地消息表 |
| MQ 存储 | 宕机丢消息 | 队列持久化 + 镜像队列 |
| 消费者接收 | 消息丢失 | 手动 ack + 失败重试 |
| 消费处理 | 重复消费 | 幂等(Redis SETNX / DB 唯一索引) |
十一、本章小结
| 要点 | 关键 |
|---|---|
| MQ 价值 | 异步、解耦、削峰 |
| 选型 | RabbitMQ(企业业务)/ RocketMQ(阿里电商)/ Kafka(日志流) |
| 五种模式 | Hello / Work / Fanout / Direct / Topic |
| 可靠发送 | ConfirmCallback + 本地消息表 |
| 可靠消费 | 手动 ack + 失败重试 |
| 幂等 | Redis SETNX / DB 唯一索引 |
| 死信 | DLX 自动收容失败消息 |
动手练习
练习 1:基础题
实现订单创建 → 通知库存扣减的简单流程:生产者发送 order.created,消费者扣减库存并打印日志。
练习 2:进阶题
为消费者加上:手动 ack + 幂等(Redis SETNX)+ 失败入死信队列 + 死信告警。模拟消费者抛异常,验证消息最终进入死信队列。
练习 3:思考题
设计一个「用户注册成功后发送欢迎邮件」的功能。用 MQ 异步发送,对比同步发送的优缺点。
下一章:第 63 章:Kafka 与流处理 →