云计算百科
云计算领域专业知识百科平台

RabbitMQ 如何保证消息不丢?Confirm、持久化与 ACK

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 提供了一系列可靠性机制,让应用能够显著降低消息丢失的风险,并在失败时进行确认、重投等处理。

真正的分布式系统还需要考虑:

网络异常
服务器故障
磁盘故障
程序异常
消息重复
消费者重试
集群故障
业务事务

赞(0)
未经允许不得转载:网硕互联帮助中心 » RabbitMQ 如何保证消息不丢?Confirm、持久化与 ACK
分享到: 更多 (0)

评论 抢沙发

评论前必须登录!