第 54 章:消息队列
学习目标
- 知道为啥用 MQ
- 掌握 RabbitMQ 发送 / 消费
- 掌握 Kafka 入门
一、为啥用 MQ
解耦 + 异步 + 削峰:
下单 → 同步:扣库存 + 发短信 + 推送,3 个调用 = 300ms
下单 → 异步:扣库存(同步) + 发消息(MQ),1 个调用 = 100ms
↓
短信服务 / 推送服务 异步消费二、RabbitMQ
2.1 集成
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest2.2 核心概念
| 概念 | 含义 |
|---|---|
| Producer | 消息生产者 |
| Consumer | 消息消费者 |
| Queue | 队列,存消息 |
| Exchange | 交换机,路由消息 |
| RoutingKey | 路由 key |
| Binding | 交换机和队列的绑定规则 |
2.3 三种交换机
| 交换机 | 路由规则 |
|---|---|
direct | RoutingKey 完全匹配 |
topic | * 匹配一个单词,# 匹配多个 |
fanout | 广播,所有绑定的队列都收 |
2.4 配置
@Configuration
public class RabbitMQConfig {
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue").build(); // 持久化
}
@Bean
public DirectExchange orderExchange() {
return new DirectExchange("order.exchange");
}
@Bean
public Binding binding(Queue orderQueue, DirectExchange orderExchange) {
return BindingBuilder.bind(orderQueue).to(orderExchange).with("order.create");
}
}💡 三个 bean 各管什么:
Queue= 信箱(消息暂存地)DirectExchange= 邮局分拣机(一个业务通常只建一个)Binding= 分拣规则("RoutingKey = X 的信 → 投到 Y 信箱")生产者和消费者互不认识:
- 🟢 生产者:
rabbitTemplate.convertAndSend("exchange 名", "RoutingKey", 内容)— 只管扔- 🔵 消费者:
@RabbitListener(queues = "队列名")— 只管收中间靠
Exchange + Binding自动路由。加新业务只需加一个 Queue + Binding,生产者代码完全不用改。
💡 简单场景可以省略 Exchange + Binding:RabbitMQ 有一个默认交换机
""(自动存在),所有队列自动绑在它上面,且路由 key = 自己的队列名。java// 只配 Queue,没配 Exchange / Binding @Bean public Queue smsQueue() { return QueueBuilder.durable("sms.queue").build(); }java// 生产者:1 个 String 参数,Spring 帮你翻译成 (默认交换机, "sms.queue" 作为路由 key) rabbitTemplate.convertAndSend("sms.queue", msg); // ✅ 找到 sms.queue但生产建议显式配 Exchange + Binding:加新业务、加路由规则都清晰,排查问题方便。只有"1 个业务 1 个队列"的玩具场景才用默认交换机。
2.5 发送
@Service
@RequiredArgsConstructor
public class OrderProducer {
private final RabbitTemplate rabbitTemplate;
public void sendOrderCreated(Order order) {
rabbitTemplate.convertAndSend(
"order.exchange", // 交换机
"order.create", // 路由 key
order // 消息体
);
}
}2.6 消费
@Component
public class OrderConsumer {
@RabbitListener(queues = "order.queue")
public void handle(Order order) {
log.info("收到订单: {}", order);
try {
// 业务逻辑
} catch (Exception e) {
log.error("消费失败", e);
throw e; // 抛出 → 进入重试 / 死信队列
}
}
}2.7 消息可靠性
💡 核心问题:消息会丢。RabbitMQ 用两层确认防丢:生产确认(生产者→broker)和消费者 ACK(broker→消费者)。默认两层都不开,消息丢了都不知道。
💡 broker 是什么? = 装在服务器上的 RabbitMQ 服务进程(
docker run rabbitmq启起来的那玩意儿),生产者不直接找消费者,而是"扔给 broker,broker 转给消费者"。broker 就是 RabbitMQ、Kafka、RocketMQ 等消息系统的统称。
@RabbitListener(queues = "order.queue")
public void handle(
Order order, // ① 消息体(自动反序列化的对象)
Channel channel, // ② AMQP 通道,用来 ACK/NACK
@Header(AmqpHeaders.DELIVERY_TAG) long tag) { // ③ 消息编号,告诉 broker 回的是哪条
try {
process(order); // ④ 真正的业务逻辑
channel.basicAck(tag, false); // ⑤ 手动 ACK:处理成功,broker 可以删消息
} catch (Exception e) {
channel.basicNack(tag, false, false); // ⑥ 处理失败:第 1 个 false=不批量,第 2 个 false=不重回队列,进死信
}
}API 拆解:
| API | 含义 |
|---|---|
basicAck(tag, false) | tag = 这条消息的编号;第二个 false = 单条确认(不批量) |
basicNack(tag, false, false) | 三个参数:①消息编号 ②是否批量 ③是否重新入队,false = 不重试,进死信队列;true = 重新入队再来一遍 |
acknowledge-mode: manual 必须配:不配这个,框架默认自动 ACK,你写的 basicAck/basicNack 就不生效。
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # ① 关键!手动 ACK 才需要写 basicAck/basicNack
prefetch: 10 # ② 一次预取 10 条消息并发处理,默认 250
retry:
enabled: true # ③ 失败时本地重试 3 次
max-attempts: 3
initial-interval: 2s⚠️ 坑 1:生产者也要开启确认(
publisher-confirm-type=correlated),否则convertAndSend一调用就返回,你以为消息送到了,其实只是字节流推到 TCP——如果 RabbitMQ 在收到前宕机,消息就丢了。💡 现实建议:绝大多数项目用默认自动 ACK + 简单重试配置就够(代码里不用写 Channel 那段),本节是严格不丢(订单、支付)的进阶方案,知道有这套即可。
三、Kafka
3.1 集成
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all # 全部副本都确认
retries: 3
consumer:
group-id: my-app
key-deserializer: ...
value-deserializer: ...
auto-offset-reset: earliest3.2 发送 / 消费
@Service
@RequiredArgsConstructor
public class OrderEventProducer {
private final KafkaTemplate<String, Order> kafkaTemplate;
public void send(Order order) {
kafkaTemplate.send("order-topic", order.getId().toString(), order);
}
}
@Component
public class OrderEventConsumer {
@KafkaListener(topics = "order-topic")
public void handle(Order order) {
log.info("收到: {}", order);
}
}3.3 Kafka vs RabbitMQ
| 场景 | 选 |
|---|---|
| 业务消息(订单、通知) | RabbitMQ |
| 日志 / 大数据流 | Kafka |
| 严格顺序 | Kafka |
| 复杂路由 | RabbitMQ |
四、本章小结
| 要点 | 关键 |
|---|---|
| RabbitMQ | 业务消息,可靠投递 |
| Kafka | 日志 / 流处理,高吞吐 |
| 可靠性 | 生产确认 + 消费手动 ACK |
| 重试 | 配置最大重试,失败进死信 |
动手练习
- 用 RabbitMQ 实现"下单后异步发短信"
- 用
@RabbitListener消费消息,故意抛异常看重试