RabbitMQ 进阶到实战:从消息可靠性到延迟队列,一篇就够了(附生产级案例解析)
系列文章:
- 🟢 入门篇:[RabbitMQ 入门到实战:从 Simple Queue 到 Topic Exchange,一篇就够了(附实战案例解析)](https://blog.csdn.net/yurenpai/article/details/161774820?fromshare=blogdetail&sharetype=blogdetail&sharerId=161774820&sharerefer=PC&sharesource=yurenpai&sharefrom=from_link)
- 🔵 进阶篇:本文 ← 你在这里
适合人群:已掌握 Simple Queue、Work Queue、三种 Exchange(Fanout/Direct/Topic)及 Spring AMQP 基本用法,准备上生产的开发者。
目录
- 一、前言:入门为什么不够用
- 二、消息可靠性(上):Producer Confirm & Return
- 三、消息可靠性(下):Consumer ACK
- 四、消息转换器:Jackson2JsonMessageConverter
- 五、死信队列(DLX):重试与异常兜底
- 六、延迟队列:TTL + DLX 与延迟插件
- 七、消息幂等性:防止重复消费
- 八、Lazy Queue:海量消息堆积的解法
- 九、集群与高可用:Quorum Queue
- 十、Spring Cloud Stream:更高级的封装(按需选用)
- 十一、总结:生产环境能力清单
一、前言:入门为什么不够用
入门阶段我们学会了"怎么发消息、怎么收消息",但真到了生产环境,问题就变了:
入门篇关心 进阶篇关心
─────────────────────────────────────────────────
消息能不能发出去? → 发了,到底有没到 Broker?
消费者能不能收到? → 收到了,处理成功了吗?
Queue 存着消息就行 → 处理失败了,消息去哪了?
Exchange 能路由就行 → 服务器挂了,消息会丢吗?
一个词概括进阶篇的核心:可靠性。
真实生产环境的六大拷问:
| 场景 | 致命问题 |
|---|---|
| Producer 发送后网络闪断 | 消息到底发成功了没?要不要重发? |
| Consumer 处理到一半宕机 | 消息算消费了还是丢了? |
| 消费失败(扣库存异常) | 重试几次?重试失败了怎么办? |
| 下单 30 分钟未支付 | 怎么自动取消订单? |
| 双十一流量洪峰 | 队列堆了几百万条,内存撑得住吗? |
| RabbitMQ 自身宕机 | 消息有没有持久化?集群有没有备份? |
本文逐一击破。

二、消息可靠性(上):Producer Confirm & Return
2.1 问题:消息到底发出去了没?
默认的 rabbitTemplate.convertAndSend() 是"即发即忘"(fire-and-forget):
rabbitTemplate.convertAndSend("order.exchange", "order.create", orderMsg);
// 这行执行完,消息真的到了吗?不知道。😕
三种静默丢消息的场景:
Producer → [网络闪断] → Exchange ❌ → 静默丢弃,不报错
Producer → Exchange ✅ → [RoutingKey 写错] → Queue ❌ → 静默丢弃,不报错
Producer → Exchange ✅ → Queue ✅ → [Consumer 宕机] → 静默丢弃,不报错
2.2 开启 Producer Confirm
spring:
rabbitmq:
publisher-confirm-type: correlated # 异步回调,推荐
publisher-returns: true # 开启路由失败通知(需配合 mandatory)
template:
mandatory: true # ⚠️ 必须开启,否则 Return 回调不生效
publisher-confirm-type 三种模式:
| 模式 | 行为 | 推荐 |
|---|---|---|
none |
关闭确认,性能最高,可靠性最低 | ❌ |
simple |
同步阻塞等待 ACK,性能差 | ❌ 已过时 |
correlated |
异步回调,不阻塞发送线程 | ✅ 生产首选 |
2.3 Confirm 回调(消息 → Exchange)
@Configuration
public class RabbitConfirmConfig {
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate template = new RabbitTemplate(connectionFactory);
// ① Confirm 回调:消息是否到达 Exchange
template.setConfirmCallback((correlationData, ack, cause) -> {
String msgId = correlationData != null ? correlationData.getId() : "unknown";
if (ack) {
log.info("✅ 消息到达 Exchange,msgId={}", msgId);
} else {
log.error("❌ 未到达 Exchange,msgId={},原因={}", msgId, cause);
// 补偿:重发 / 入库 / 告警
}
});
// ② Return 回调:Exchange → Queue 路由失败
template.setMandatory(true); // ⚠️ 自定义 Bean 时必须在代码里设置,YAML 不生效
template.setReturnsCallback(returned -> {
log.error("❌ 路由失败!exchange={},routingKey={},replyCode={},body={}",
returned.getExchange(),
returned.getRoutingKey(),
returned.getReplyCode(),
new String(returned.getMessage().getBody()));
// 补偿:修正 RoutingKey / 人工介入
});
return template;
}
}
发送时带上 CorrelationData(回调时靠它识别是哪条消息):
CorrelationData data = new CorrelationData(UUID.randomUUID().toString());
rabbitTemplate.convertAndSend("order.exchange", "order.create", orderMsg, data);
2.4 两种回调的分工
| 回调 | 监听什么 | 触发条件 |
|---|---|---|
| ConfirmCallback | Producer → Exchange | Exchange 不存在 / 权限不足 |
| ReturnsCallback | Exchange → Queue | RoutingKey 匹配不到任何队列 |
关键:
mandatory: true必须开启,否则消息路由失败时直接丢弃,不会触发 ReturnsCallback。
2.5 完整流程
┌─ ack=true → ✅ 到达 Exchange,继续路由到 Queue
Producer ── 发消息 ──┤
│ ┌─ 路由成功 → Queue 等待消费
│ Exchange ──┤
│ └─ 路由失败 → 触发 ReturnsCallback(需 mandatory=true)
│
└─ ack=false → ❌ 未到达 Exchange
↓
触发 ConfirmCallback → 重发/入库/告警
2.6 补充:消息持久化投递
Confirm + Return 解决的是"发没发到"的问题,但消息本身如果只存在内存,Broker 重启一样丢。生产者发送时可以显式标记持久化:
// Spring Boot 自动配置了 ObjectMapper,直接注入即可
@Autowired
private ObjectMapper objectMapper;
public void sendOrder(Order order) {
Message message = MessageBuilder
.withBody(objectMapper.writeValueAsBytes(order))
.setDeliveryMode(MessageDeliveryMode.PERSISTENT) // 👈 消息持久化到磁盘
.setMessageId(UUID.randomUUID().toString()) // 消息唯一ID(幂等用)
.build();
CorrelationData data = new CorrelationData(message.getMessageProperties().getMessageId());
rabbitTemplate.convertAndSend("order.exchange", "order.create", message, data);
}
三重持久化保障:
Exchange Durable+Queue Durable+Message Persistent,缺一不可。前两者在声明时默认 durable=true,第三者需发送时显式指定。
三、消息可靠性(下):Consumer ACK
3.1 问题:消费者真的处理成功了吗?
默认是自动 ACK:消息一到 Consumer,RabbitMQ 立即从队列删除。
消息到达 Consumer → RabbitMQ 立即删除 → Consumer 处理到一半崩了 💥
↓
消息永久丢失 😱
3.2 手动 ACK
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # 手动确认
prefetch: 1 # 每次只拉取1条;生产可按需调大(如250)提升吞吐
@Component
public class OrderConsumer {
@RabbitListener(queues = "order.queue")
public void handle(Message message, Channel channel) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
// 1. 业务处理
processOrder(new String(message.getBody()));
// 2. 成功 → ACK
channel.basicAck(deliveryTag, false);
log.info("✅ 处理成功,deliveryTag={}", deliveryTag);
} catch (Exception e) {
log.error("❌ 处理失败", e);
// 3. 失败 → NACK + 不重新入队(给 DLX 兜底)
channel.basicNack(deliveryTag, false, false);
// ↑
// false = 不 requeue,进入死信队列
}
}
}
3.3 ACK / NACK / Reject 对比
| 方法 | 效果 | 使用场景 |
|---|---|---|
basicAck |
确认删除 | 业务处理成功 |
basicNack(requeue=true) |
拒收,重新入队 | 临时故障,重试可能恢复 |
basicNack(requeue=false) |
拒收,不入队 | 永久错误,丢入 DLX |
basicReject |
同 Nack(false),仅单条 | 单条拒绝 |
3.4 手动 ACK 三大坑
⚠️ 忘记 ACK → 消息一直 Unacked → 消费者断开后重新入队 → 重复消费
⚠️ 忘记 NACK → 消息永远悬在 Unacked → 内存泄漏
⚠️ NACK + requeue=true 死循环 → 无限重试 → CPU 打满
✅ 正确姿势:NACK + requeue=false + 死信队列兜底(见第五章)
四、消息转换器:Jackson2JsonMessageConverter
4.1 默认 JDK 序列化的坑
RabbitTemplate 默认用 SimpleMessageConverter(JDK 序列化),后果:
发送 Order 对象
↓
JDK 序列化 → 二进制乱码
↓
RabbitMQ 控制台:□□□□□□□□(完全不可读)
↓
非 Java 消费者看不懂 / 不同版本 JDK 反序列化失败
4.2 切到 JSON(一行配置)
@Configuration
public class RabbitConfig {
@Bean
public MessageConverter messageConverter() {
return new Jackson2JsonMessageConverter(); // 👈 就这一行
}
}
消费者两种接收方式:
// 方式一:收 String,自己用 Jackson 反序列化(推荐,简单可控)
@RabbitListener(queues = "order.queue")
public void handle(String json) {
Order order = objectMapper.readValue(json, Order.class);
}
// 方式二:直接收对象(需消息头带 _TypeId_,Spring 自动反序列化)
@RabbitListener(queues = "order.queue")
public void handle(Order order) {
// Spring 利用 Jackson 自动将 JSON → Order 对象
}
4.3 对比
| 维度 | JDK 序列化 | Jackson JSON |
|---|---|---|
| 可读性 | ❌ 二进制乱码 | ✅ 明文 JSON |
| 跨语言 | ❌ 仅 Java | ✅ 任意语言 |
| 管理界面查看 | ❌ 不可读 | ✅ 直接可读 |
| 推荐度 | ❌ 永远别用 | ✅ 生产铁律 |
生产铁律:引入 RabbitMQ 第一步就换掉默认序列化,用
Jackson2JsonMessageConverter。
五、死信队列(DLX):重试与异常兜底
5.1 什么情况消息会变成"死信"?
| 情况 | 说明 |
|---|---|
| 被拒绝 | Consumer 执行 basicNack/basicReject 且 requeue=false |
| 过期 | 消息 TTL 到期未被消费 |
| 队列满了 | 队列达到最大长度,溢出消息被"挤"出去 |

5.2 DLX 模型
┌─ 处理成功 → ACK ✅
订单队列(正常队列) ──┤
└─ 处理失败 → NACK(requeue=false) → 死信交换机(DLX)
↓
死信队列(DLQ)
↓
死信消费者
(记录DB / 告警 / 人工修复)
5.3 完整代码
声明正常队列 + 死信交换机 + 死信队列:
@Configuration
public class DLXConfig {
// ========== 死信交换机 & 死信队列 ==========
@Bean
public DirectExchange dlxExchange() {
return new DirectExchange("dlx.exchange");
}
@Bean
public Queue dlq() {
return QueueBuilder.durable("dlq").build();
}
@Bean
public Binding dlxBinding() {
return BindingBuilder.bind(dlq()).to(dlxExchange()).with("dlx.routing.key");
}
// ========== 正常业务队列(绑定 DLX)==========
@Bean
public Queue orderQueue() {
return QueueBuilder
.durable("order.queue")
.deadLetterExchange("dlx.exchange") // 👈 核心:绑定死信交换机
.deadLetterRoutingKey("dlx.routing.key")
.ttl(10000) // 可选:消息 10 秒过期
.maxLength(1000) // 可选:队列最大长度
.build();
}
@Bean
public DirectExchange orderExchange() {
return new DirectExchange("order.exchange");
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue()).to(orderExchange()).with("order.create");
}
}
死信消费者:兜底处理:
@Component
public class DeadLetterConsumer {
@RabbitListener(queues = "dlq")
public void handleDeadLetter(Message message) {
String body = new String(message.getBody());
// 死信原因等元信息在 header 里
List<Map<String, Object>> deathInfo = (List<Map<String, Object>>)
message.getMessageProperties().getHeaders().get("x-death");
log.warn("⚠️ 收到死信:body={},deathInfo={}", body, deathInfo);
// 策略:
// 1. 持久化到 DB,方便人工排查
// 2. 发送钉钉/企微告警
// 3. 特定类型修复后可手动重新投递
}
}
5.4 重试策略:逐级重试 → 死信兜底
第1次失败 → NACK(requeue=true) → 重新入队
第2次失败 → NACK(requeue=true) → 重新入队
第3次失败 → NACK(requeue=true) → 重新入队
第4次失败 → NACK(requeue=false) → 进入死信队列 → 人工介入 🛎️
实现方式:
方案A:消息头记录 x-retry-count,消费者自己判断
方案B:Spring Retry 配置(更简单,见下方 YAML)
# 方式B:Spring RabbitMQ 自带重试
spring:
rabbitmq:
listener:
simple:
retry:
enabled: true
max-attempts: 3
initial-interval: 1000ms
multiplier: 2.0 # 间隔递增:1s → 2s(初次→1次重试→2次重试)
六、延迟队列:TTL + DLX 与延迟插件
6.1 典型场景
下单 30 分钟未支付 → 自动取消订单,释放库存
用户注册 24 小时未激活 → 清理账号
会议前 15 分钟 → 发送提醒通知
不用定时任务扫表(SELECT ... WHERE create_time < now() - 30min),延迟队列更优雅。
6.2 两种实现方式
┌─ 方案一:消息 TTL + DLX(原生支持,无需插件)
│ 消息设 TTL → 过期变死信 → DLX 路由到处理队列
│ 缺点:消息只有在队列头部才会被检测到过期(头阻塞)
│
└─ 方案二:rabbitmq_delayed_message_exchange 插件(⭐ 推荐)
安装后 Exchange 原生支持延迟投递,不受头阻塞影响

6.3 方案一:TTL + DLX
@Configuration
public class DelayQueueConfig {
// ========== 实际处理队列(接收过期消息)==========
@Bean
public Queue cancelOrderQueue() {
return QueueBuilder.durable("order.cancel.queue").build();
}
@Bean
public DirectExchange cancelOrderExchange() {
return new DirectExchange("order.cancel.exchange");
}
@Bean
public Binding cancelOrderBinding() {
return BindingBuilder
.bind(cancelOrderQueue())
.to(cancelOrderExchange())
.with("order.cancel");
}
// ========== 延迟队列(设 TTL,到期自动进死信 → 处理队列)==========
@Bean
public Queue delayQueue30min() {
return QueueBuilder
.durable("delay.30min.queue")
.deadLetterExchange("order.cancel.exchange") // "死信" = 实际处理交换机
.deadLetterRoutingKey("order.cancel")
.ttl(30 * 60 * 1000) // 30 分钟
.build();
}
@Bean
public DirectExchange delayExchange() {
return new DirectExchange("delay.exchange");
}
@Bean
public Binding delayBinding30min() {
return BindingBuilder
.bind(delayQueue30min())
.to(delayExchange())
.with("delay.30min");
}
}
// 发送延迟消息:投到延迟队列,30 分钟后自动进入处理队列
rabbitTemplate.convertAndSend("delay.exchange", "delay.30min", orderMsg);
6.4 方案二:延迟交换机插件(⭐ 生产推荐)
# 1. 下载插件 .ez 文件(以 RabbitMQ 3.12 为例,版本号需与你的 RabbitMQ 匹配)
# 下载地址:https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases
docker exec -it rabbitmq bash -c \
"wget -P /plugins \
https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/\
releases/download/v3.12.0/rabbitmq_delayed_message_exchange-3.12.0.ez"
# 2. 启用插件
docker exec -it rabbitmq rabbitmq-plugins enable rabbitmq_delayed_message_exchange
@Bean
public CustomExchange delayedExchange() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct");
return new CustomExchange(
"order.delayed.exchange",
"x-delayed-message", // 延迟交换机类型
true, false, args
);
}
发送消息时通过 MessagePostProcessor 设置延迟时间:
// 发送时指定延迟时间(每条消息可不同延迟)
rabbitTemplate.convertAndSend(
"order.delayed.exchange", "order.cancel", orderMsg,
msg -> {
msg.getMessageProperties().setDelay(30 * 60 * 1000); // 30 分钟
return msg;
}
);
6.5 两种方案对比
| 维度 | TTL + DLX | 延迟插件 |
|---|---|---|
| 队列数量 | 每个延迟时间一个队列 | 一个交换机搞定 |
| 消息顺序 | 受队列头部阻塞影响 | ✅ 不受影响 |
| 灵活性 | ❌ 队列 TTL 固定 | ✅ 每条消息可不同延迟 |
| 安装 | 无需插件 | 需安装插件 |
| 生产推荐 | 简单场景 | ⭐ 首选 |
七、消息幂等性:防止重复消费
7.1 重复消费从哪来?
场景1:Consumer 处理成功 → 发 ACK → 网络闪断 ACK 丢失 → MQ 重新投递
场景2:Producer Confirm 超时 → 重发 → 同一条消息发了两遍
场景3:NACK requeue → 重新消费同一条消息
RabbitMQ 保证的是 at-least-once(至少一次),不是 exactly-once。消费者必须自己做幂等。
7.2 四种幂等方案
| 方案 | 实现 | 适用场景 |
|---|---|---|
| DB 唯一索引 | 消息 ID 建唯一索引,INSERT IGNORE | 通用,最靠谱 |
| Redis SETNX | 消息 ID 做 key,set 成功才处理 | 高性能场景 |
| 业务状态机 | 订单已支付 → 跳过支付操作 | 有明确状态流转 |
| 版本号 | 消息带 version,小于当前版本跳过 | 更新类操作 |
7.3 DB 唯一索引方案(实战)
CREATE TABLE mq_consume_record (
msg_id VARCHAR(64) PRIMARY KEY, -- 消息唯一 ID
create_time DATETIME DEFAULT NOW()
);
@Component
public class IdempotentConsumer {
@RabbitListener(queues = "order.queue")
public void handle(Message message, Channel channel) throws IOException {
String msgId = message.getMessageProperties().getMessageId();
long deliveryTag = message.getMessageProperties().getDeliveryTag();
// 1. 幂等检查:利用唯一索引防重
try {
mqConsumeRecordMapper.insert(new MqConsumeRecord(msgId));
} catch (DuplicateKeyException e) {
log.warn("重复消息,跳过:{}", msgId);
channel.basicAck(deliveryTag, false); // 直接 ACK
return;
}
// 2. 执行业务
try {
doBusiness(message);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
// 业务失败,删除消费记录允许重试
mqConsumeRecordMapper.delete(msgId);
channel.basicNack(deliveryTag, false, true);
}
}
}
7.4 Redis 方案(高性能)
String msgId = message.getMessageProperties().getMessageId();
Boolean locked = redisTemplate.opsForValue()
.setIfAbsent("mq:consumed:" + msgId, "1", Duration.ofHours(24));
if (Boolean.FALSE.equals(locked)) {
channel.basicAck(deliveryTag, false); // 已消费,跳过
return;
}
生产铁律:Consumer 必须做幂等,不要指望 MQ 不重复投递。这是血的教训。
八、Lazy Queue:海量消息堆积的解法
8.1 默认队列的内存困境
双十一流量洪峰
↓
100 万条消息堆积在内存
↓
内存爆炸 → OOM → RabbitMQ 崩溃 💥
8.2 Lazy Queue 原理
普通队列:消息先存内存 → 内存满了才刷盘
Lazy Queue:消息直接写磁盘 → 内存只存索引 → 支持 TB 级堆积
8.3 声明
// 方式1:QueueBuilder
@Bean
public Queue lazyQueue() {
return QueueBuilder.durable("lazy.queue")
.lazy() // 👈 就这一行
.build();
}
// 方式2:注解
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "lazy.queue",
arguments = @Argument(name = "x-queue-mode", value = "lazy")),
exchange = @Exchange(name = "lazy.exchange", type = ExchangeTypes.TOPIC),
key = "lazy.#"
))
public void handle(String msg) { }
8.4 默认队列 vs Lazy Queue
| 维度 | 默认队列 | Lazy Queue |
|---|---|---|
| 存储位置 | 内存优先 | 磁盘优先 |
| 堆积能力 | 受内存限制(GB 级) | 受磁盘限制(TB 级) |
| 消费速度 | 快(内存读取) | 略慢(磁盘 I/O) |
| 适用场景 | 日常低堆积 | 高吞吐、大堆积、秒杀 |
建议:不需要所有队列都开 Lazy。日常流量用默认队列(性能好),已知会大堆积的(秒杀、日志收集)才用 Lazy Queue。
九、集群与高可用:Quorum Queue
9.1 单机的致命问题
RabbitMQ 宕机 → 所有消息不可用 → 系统瘫痪 💀
9.2 镜像队列(传统方案,了解即可)
┌── Node 1 (Master) ── 所有读写
│ │
集群 ───┼── Node 2 (Mirror) ── 自动同步
│ │
└── Node 3 (Mirror) ── 自动同步
Master 挂了 → Mirror 自动提升为新 Master
问题:最终一致性,故障时可能丢消息
9.3 Quorum Queue(RabbitMQ 3.8+,⭐ 新项目首选)
基于 Raft 共识协议,强一致性,自动故障转移。
@Bean
public Queue quorumOrderQueue() {
return QueueBuilder.durable("order.queue")
.quorum() // 👈 声明为 Quorum Queue
.build();
}
// 注解方式
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "order.queue",
arguments = @Argument(name = "x-queue-type", value = "quorum")),
exchange = @Exchange(name = "order.exchange", type = ExchangeTypes.TOPIC),
key = "order.#"
))
9.4 镜像队列 vs Quorum Queue
| 维度 | 镜像队列 | Quorum Queue |
|---|---|---|
| 一致性 | 最终一致(可能丢消息) | 强一致(Raft) |
| 故障转移 | 手动或半自动 | 全自动 |
| 性能 | 较高 | 略低(Raft 开销) |
| 故障时消息重复 | 可能重复 | 不重复 |
| 消息优先级 | ✅ 支持 | ❌ 不支持 |
| 推荐度 | 老项目在用 | ⭐ 新项目首选 |

建议:RabbitMQ 3.8+ 新项目直接用 Quorum Queue,老项目逐步从镜像队列迁移。
十、Spring Cloud Stream:更高级的封装(按需选用)
⚠️ API 版本说明:下面展示的是 Spring Cloud Stream 2.x 的
@EnableBinding/@StreamListener注解风格。自 Spring Cloud Stream 3.1+ 起,这些注解已废弃,推荐使用函数式编程模型(java.util.function.Consumer/Function/SupplierBean)。此处仍展示旧 API 是为了帮助理解核心概念,并兼容存量的老项目代码。新项目请直接使用函数式模型。
10.1 为什么要再封装一层?
原生 Spring AMQP:
rabbitTemplate.convertAndSend("order.exchange", "order.key", msg)
@RabbitListener(queues = "order.queue")
↓
强耦合 RabbitMQ API
↓
如果以后切到 Kafka / RocketMQ?全部代码要改 😰
Spring Cloud Stream 的思路:业务代码只面对 Binder 抽象,底层换 MQ 只改配置。
10.2 核心概念
Source(消息源) → 发送消息
Sink(消息目的地) → 接收消息
Binder → 绑定到具体 MQ(RabbitMQ / Kafka)
Channel → 消息通道
10.3 生产者示例
// 定义 Binding
public interface OrderSource {
String OUTPUT = "order-output";
@Output(OUTPUT)
MessageChannel output();
}
// 启动类
@SpringBootApplication
@EnableBinding(OrderSource.class)
public class Application { }
// 发送
@Autowired
private OrderSource orderSource;
public void send(Order order) {
orderSource.output().send(
MessageBuilder.withPayload(objectMapper.writeValueAsString(order)).build()
);
}
10.4 消费者示例
public interface OrderSink {
String INPUT = "order-input";
@Input(INPUT)
SubscribableChannel input();
}
@Component
public class OrderHandler {
@StreamListener(OrderSink.INPUT)
public void handle(String msg) {
log.info("收到订单:{}", msg);
}
}
10.5 切 Kafka 只改这里
依赖换一下:
spring-cloud-starter-stream-rabbit
↓
spring-cloud-starter-stream-kafka
yaml 改一下 Broker 地址 → 代码零改动。
10.6 务实建议
| 场景 | 建议 |
|---|---|
| 单一 MQ,不会换 | 直接用 Spring AMQP,别多套一层 |
| 可能切换 MQ 或多 MQ 并存 | 用 Spring Cloud Stream |
| 团队大,需要统一规范 | 用 Stream 约束 API |
个人建议:大部分公司不会频繁换 MQ,直接用 Spring AMQP 就够。多一层抽象就多一层排查成本,不建议为"未来可能切"而过度设计。
十一、总结:生产环境 RabbitMQ 能力清单
11.1 完整能力进阶
入门 进阶
──────── ──────────
发消息 Simple Queue → Producer Confirm ✅
收消息 @RabbitListener → Consumer Manual ACK ✅
消息格式 JDK 序列化(坑) → Jackson2Json ✅
路由 Fanout/Direct/Topic → 延迟交换机
异常处理 丢了就丢了 → DLX + 重试策略 ✅
重复消费 不管 → 幂等(DB/Redis) ✅
大流量 默认队列 → Lazy Queue
高可用 单机 → Quorum Queue 集群 ✅
框架封装 Spring AMQP → Spring Cloud Stream(按需)
标 ✅ = 上生产前必须搞定。
11.2 生产环境推荐配置
spring:
rabbitmq:
host: ${RABBIT_HOST:localhost}
port: ${RABBIT_PORT:5672}
username: ${RABBIT_USER:admin}
password: ${RABBIT_PASS:admin}
# ✅ 生产者确认
publisher-confirm-type: correlated
publisher-returns: true
template:
mandatory: true # 必须配合 ReturnCallback
# ✅ 消费者确认 + 重试
listener:
simple:
acknowledge-mode: manual
prefetch: 250
retry:
enabled: true
max-attempts: 3
initial-interval: 1000ms
multiplier: 2.0
11.3 生产完整链路(一条消息的一生)
Producer
│ 带 CorrelationData + messageId
▼
Publisher Confirm(到没到 Exchange?)
│
▼
Durable Exchange(持久化)
│
▼
Durable / Quorum / Lazy Queue(持久化 + 高可用 + 抗堆积)
│
▼
Consumer Manual ACK(成功才确认)
│
├── ✅ 成功 → ACK → 幂等校验通过 → 业务完成
│
└── ❌ 失败 → Retry 3次 → 仍失败 → DLX 死信队列 → 告警/人工处理
11.4 学习路线
🟢 入门(已掌握) 🟡 进阶(本文) 🟠 高级(进阶方向)
─────────────────── ────────────────── ───────────────
Simple Queue Producer Confirm 消息轨迹追踪
Work Queue Consumer Manual ACK Federation 联邦
Fanout Exchange Jackson2Json 转换器 Shovel 跨机房同步
Direct Exchange DLX 死信队列 Prometheus 监控告警
Topic Exchange TTL + 延迟队列 性能调优
Spring AMQP 集成 Lazy Queue 源码阅读
@RabbitListener Quorum Queue 集群
Spring Cloud Stream
你已经在这 ✅ 你现在在这 ✅ 下一站 🚀
系列文章导航:
| 阶段 | 标题 | 关键词 |
|---|---|---|
| 🟢 入门 | [RabbitMQ 入门到实战:从 Simple Queue 到 Topic Exchange,一篇就够了(附实战案例解析)](TODO_替换为上一篇CSDN链接) | Simple Queue、Work Queue、Fanout/Direct/Topic Exchange、Spring AMQP |
| 🔵 进阶 | 本文 | Producer Confirm、Consumer ACK、DLX、延迟队列、幂等性、Lazy Queue、Quorum Queue |
| 🟠 高级 | 敬请期待 | 消息轨迹追踪、Federation 联邦、Shovel 跨机房、Prometheus 监控告警 |
如果本文对你有帮助,欢迎点赞、收藏、关注!
更多推荐

所有评论(0)