RabbitMQ 最基本的消息链路:
Producer
↓
Exchange
↓
Queue
↓
Consumer
看起来很简单,但真正用于业务以后,一个非常重要的问题就出现了:
如果中间某个环节出现异常,消息会不会丢?
例如:
Producer
↓
RabbitMQ
生产者发送消息时网络突然断开怎么办?
或者:
RabbitMQ
↓
Consumer
消费者拿到消息后,业务还没执行完程序就崩了怎么办?
所以 RabbitMQ 的可靠性并不是依靠某一个功能完成的,而是需要从整条消息链路考虑。
可以简单分成三段:
Producer
│
│ Publisher Confirm
↓
RabbitMQ
│
│ 持久化
↓
Queue
│
│ Consumer ACK
↓
Consumer
本文就围绕这三个机制展开。
1. 消息到底可能在哪里丢?
先假设订单服务发送一条消息:
订单创建成功
完整流程:
订单服务
↓
RabbitMQ
↓
订单队列
↓
短信服务
这里至少存在三个风险。
第一种:生产者发送失败
Producer
↓
× 网络异常
↓
RabbitMQ
生产者执行了发送代码:
rabbitTemplate.convertAndSend(...);
但是这并不代表 RabbitMQ 一定已经成功接收到消息。
网络可能在传输过程中发生异常。
RabbitMQ 官方也明确指出,仅仅把数据写入 Socket,并不能证明消息已经成功到达并被 Broker 接收,因此需要 Publisher Confirm。
第二种:RabbitMQ 收到消息后宕机
即使 RabbitMQ 已经收到消息:
Producer
↓
RabbitMQ
↓
Queue
如果这些数据没有正确持久化,RabbitMQ 重启后消息仍然可能消失。
第三种:消费者处理到一半崩溃
例如消费者收到订单消息:
@RabbitListener(queues = "order.queue")
public void receive(String message) {
saveData();
sendSms();
}
假设刚拿到消息:
RabbitMQ
↓
Consumer
消费者突然宕机。
如果 RabbitMQ 已经认为:
这条消息消费成功了
那么消息可能已经从队列中移除。
但实际上消费者根本没有执行完业务。
因此还需要:
ACK
也就是消费者确认机制。
2. Publisher Confirm:消息真的到 RabbitMQ 了吗?
先解决第一段:
Producer → RabbitMQ
Publisher Confirm 可以理解成 RabbitMQ 给生产者的一张“回执”。
原本:
Producer
│
│ 发送消息
↓
RabbitMQ
生产者发送以后不知道结果。
开启 Confirm 后:
Producer
│
│ message
↓
RabbitMQ
│
│ ACK / NACK
↓
Producer
RabbitMQ 会告诉生产者:
消息我已经处理了
或者:
这条消息出现了问题
官方将 Publisher Confirm 定义为生产者跟踪哪些消息已经被 RabbitMQ 成功接收处理的机制。
3. Spring Boot 开启 Publisher Confirm
可以在配置文件中开启:
spring:
rabbitmq:
publisher-confirm-type: correlated
correlated 模式可以结合:
CorrelationData
让我们知道具体是哪条消息得到了确认。
发送消息:
CorrelationData correlationData =
new CorrelationData("order-10001");
rabbitTemplate.convertAndSend(
"order.exchange",
"order.created",
"订单创建成功",
correlationData
);
这里:
order-10001
相当于给消息增加了一个业务上的追踪标识。
这样发生异常以后,就可以知道:
到底是哪条消息没有确认?
4. Confirm 成功就代表消费者收到消息了吗?
不是。
这一点非常重要。
Confirm 解决的是:
Producer
↓
RabbitMQ
它能够告诉生产者:
RabbitMQ 是否已经接收处理这条消息。
但并不能代表:
Consumer 已经消费成功
所以不要把:
Publisher Confirm
理解成:
整个消息业务执行成功
它只解决生产者到 Broker 这一段的问题。
5. 消息到了 Exchange,却没有进入 Queue 怎么办?
这里还有一种比较特殊的情况。
假设发送:
rabbitTemplate.convertAndSend(
"order.exchange",
"abc",
"订单创建成功"
);
但是 Exchange 中根本没有匹配:
abc
的 Binding。
那么:
Producer
↓
Exchange
↓
×
Queue
消息虽然到达了 Exchange,却没有被路由到任何 Queue。
这时候可以使用 RabbitMQ 的 Return 机制。
Spring Boot 中可以配置:
spring:
rabbitmq:
publisher-returns: true
template:
mandatory: true
这样当消息无法路由到 Queue 时,就可以进行相应处理。
因此 Producer 端实际上需要关心两个问题:
消息有没有到 RabbitMQ?
↓
Publisher Confirm
消息有没有成功路由到 Queue?
↓
Return
这两个概念不要混在一起。
6. 持久化:RabbitMQ 重启后消息还在吗?
现在解决第二个问题:
RabbitMQ 自己挂了怎么办?
RabbitMQ 中的持久化至少需要关注:
Queue 持久化
Message 持久化
7. Queue 持久化
之前创建 Queue:
@Bean
public Queue orderQueue() {
return new Queue("order.queue");
}
Spring AMQP 这种常见声明默认就是 durable Queue。
也可以明确写:
return QueueBuilder
.durable("order.queue")
.build();
durable 表示这个 Queue 的元数据可以在 RabbitMQ 重启后恢复。
RabbitMQ 官方文档指出,durable Queue 会在 Broker 重启后恢复,而 transient Queue 则不会。
但这里有一个很容易产生的误区:
Queue 持久化,不代表 Queue 里面所有消息一定都会持久化。
8. Message 也需要持久化
RabbitMQ 对消息本身也区分:
persistent
transient
如果业务要求 RabbitMQ 重启以后消息仍然恢复,需要:
Durable Queue
+
Persistent Message
也就是:
队列要持久化
消息也要持久化
RabbitMQ 官方同样强调,对于要求数据可靠保存的场景,需要同时使用 durable Queue,并将消息发布为 persistent。
可以简单记成:
只持久化 Queue
↓
只是“箱子”还在
Queue + Message 都持久化
↓
箱子和里面的重要东西都尽量保存下来
9. 持久化是不是绝对不会丢消息?
也不是。
比如消息刚到 RabbitMQ:
RabbitMQ 收到消息
↓
准备持久化
↓
服务器突然掉电
如果生产者还不知道 RabbitMQ 是否真正安全处理了这条消息,仍然存在不确定性。
所以实际可靠性通常不是只依赖:
消息持久化
而是结合:
Publisher Confirm
+
Queue 持久化
+
Message 持久化
一起使用。
10. Consumer ACK:消费者真的处理成功了吗?
现在来看消息链路的最后一段:
RabbitMQ
↓
Consumer
消费者拿到消息并不代表:
业务处理成功
例如:
public void receive(String message) {
updateDatabase();
int a = 1 / 0;
sendSms();
}
消费者确实拿到了消息。
但是执行到一半:
发生异常
这时候 RabbitMQ 应该知道:
这条消息到底成功了没有?
因此就有了:
Consumer ACK
11. ACK 的核心思想
ACK 就是:
Acknowledgement
也就是确认。
可以理解成消费者处理完以后告诉 RabbitMQ:
Consumer:这条消息我处理成功了
↓
RabbitMQ:好,那我可以处理掉这条已确认消息了
流程:
RabbitMQ
│
│ message
↓
Consumer
│
│ 执行业务
↓
业务成功
│
│ ACK
↓
RabbitMQ
RabbitMQ 支持自动确认和显式确认模式;在需要消费者确认的情况下,未确认消息在消费者连接或 Channel 出现异常时可以重新投递。
12. Spring AMQP 手动 ACK
可以直接给监听器指定:
ackMode = "MANUAL"
例如:
@RabbitListener(
queues = "order.queue",
ackMode = "MANUAL"
)
public void receive(
String message,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag)
throws IOException {
try {
System.out.println("处理消息:" + message);
// 执行业务代码
channel.basicAck(
deliveryTag,
false
);
} catch (Exception e) {
channel.basicNack(
deliveryTag,
false,
true
);
}
}
Spring AMQP 官方也支持在 @RabbitListener 中通过 ackMode = "MANUAL" 使用 Channel.basicAck() 手动确认消息。
13. basicAck 是什么意思?
channel.basicAck(
deliveryTag,
false
);
其中:
deliveryTag
可以理解成当前消息在 Channel 中的投递编号。
而:
false
表示:
只确认当前这条消息
所以业务执行成功以后:
Consumer
↓
basicAck()
↓
RabbitMQ
告诉 RabbitMQ:
这条消息已经处理成功
14. 处理失败怎么办?
如果消费者处理失败:
channel.basicNack(
deliveryTag,
false,
true
);
最后一个参数:
true
表示:
重新入队
消息可能再次进入队列等待消费:
Consumer
↓
处理失败
↓
NACK
↓
Queue
↓
重新消费
看起来很好,但这里马上又会产生一个新问题。
假设代码本身就是错的:
int a = 1 / 0;
那么:
消费
↓
失败
↓
重新入队
↓
再次消费
↓
再次失败
↓
再次入队
就可能形成:
死循环
所以实际项目通常还需要:
重试次数限制
死信队列
异常消息处理
这也是后面需要继续学习的内容。
15. ACK 会带来另一个问题:重复消费
假设:
Consumer
↓
数据库已经修改成功
↓
准备发送 ACK
↓
Consumer 突然宕机
RabbitMQ 没有收到 ACK。
它可能认为:
这条消息还没有处理成功
于是重新投递。
结果就变成:
第一次:
余额 + 100
第二次重新消费:
余额又 + 100
也就是说:
可靠消息机制在减少消息丢失的同时,也可能带来重复消费。
所以 RabbitMQ 可靠性问题通常还会继续延伸到一个非常重要的概念:
幂等性
消费者应该尽量保证:
同一条消息执行一次
和
同一条消息执行多次
最终业务结果相同
例如可以通过:
消息唯一 ID
数据库唯一约束
Redis 去重
业务状态判断
避免重复消费造成错误。
16. 把整条可靠性链路串起来
现在再看一遍:
Producer
│
│ Publisher Confirm
↓
Exchange
│
│ Return
↓
Queue
│
│ Queue + Message 持久化
↓
RabbitMQ
│
│ Consumer ACK
↓
Consumer
每一个机制解决的问题都不一样。
Publisher Confirm
解决:
消息到底有没有成功到达 RabbitMQ?
Return
解决:
Exchange 收到了消息,
但是消息有没有路由到 Queue?
持久化
解决:
RabbitMQ 重启以后,
重要的 Queue 和消息还能不能恢复?
Consumer ACK
解决:
Consumer 到底有没有真正处理完消息?
17. RabbitMQ 能做到 100% 永远不丢消息吗?
学习 RabbitMQ 时,很容易听到:
RabbitMQ 可以保证消息不丢。
这句话其实需要更严谨一些。
更加准确地说:
RabbitMQ 提供了一系列可靠性机制,让应用能够显著降低消息丢失的风险,并在失败时进行确认、重投等处理。
真正的分布式系统还需要考虑:
网络异常
服务器故障
磁盘故障
程序异常
消息重复
消费者重试
集群故障
业务事务
网硕互联帮助中心



评论前必须登录!
注册