面试官必问:RabbitMQ 如何保证消息可靠传递?
目录
前言
在分布式系统中,消息队列作为核心组件,承担着系统解耦、流量削峰的重任。而RabbitMQ作为业界广泛使用的消息中间件,其可靠性问题一直是开发者关注的焦点。
“消息会不会丢失?”“服务器宕机怎么办?”“如何确保消息一定能到达消费者?”这些问题在日常开发和面试中都频繁出现。本文将深入探讨RabbitMQ的消息可靠传输机制,从生产者、RabbitMQ服务端到消费者三个维度,全面解析如何构建一个高可靠的消息系统。
一、消息丢失的三个环节
1. 消息的生命周期
一条消息从产生到被消费,需要经历三个关键阶段:

每个阶段都可能成为消息丢失的隐患点。
2. 丢失消息三个环节
消息丢失大概分为三个环节:
1. 生产者发送环节
因为应用程序故障,网络抖动等各种原因,生产者没有成功向RabbitMQ发送消息;
常见原因:
-
网络问题:消息在传输过程中丢失,未能到达RabbitMQ服务端;
-
路由错误:消息发送到交换机后,由于路由键配置错误,无法路由到任何队列;
-
服务端拒绝:RabbitMQ服务端内存或磁盘空间不足,拒绝接收新消息;
2. RabbitMQ处理环节
RabbitMQ自身问题,生产者成功发送给了RabbitMQ,但是RabbitMQ自身没有把消息保存好,导致消息丢失;
常见原因:
-
服务器宕机:消息暂存在内存中,服务器突然宕机导致消息丢失;
-
磁盘损坏:持久化消息已写入磁盘,但磁盘物理损坏导致数据不可恢复;
-
节点故障:集群模式下单个节点故障,该节点上的队列不可用;
3. 消费者消费环节
RabbitMQ成功将消息发送给了消费者,但是消费者在消费消息的过程中没有处理好,导致RabbitMQ将消费的消息从队列中删除;
常见原因:
-
处理异常:消费者收到消息后业务处理抛出异常,消息未正确处理;
-
自动确认陷阱:使用自动ACK,消费者接收消息后立即确认,但处理过程中宕机,消息丢失;
-
消费者宕机:消息已投递但消费者尚未处理完成,连接断开导致消息重新入队或丢失;
3. 消息可靠投递的设计思路
要保证消息不丢失,必须采取全链路覆盖的设计思路:每个环节都有对应的保障措施,形成完整的可靠性闭环。
核心设计原则:
-
确认机制:每个环节都需要确认,发送方知道消息是否成功
-
持久化:关键数据写入磁盘,防止进程或服务器重启导致丢失
-
冗余备份:集群多副本存储,应对单点故障
-
重试机制:失败后提供重试机会,不轻易放弃
-
兜底策略:最终无法处理的消息进入死信队列,人工介入
二、发布确认模式
1. Confirm确认模式
工作原理:生产者将信道设置为confirm模式,所有在该信道发布的消息都会被分配一个唯一ID。消息到达交换机后,RabbitMQ会异步发送确认给生产者。
Confirm模式有三种实现方式,分别为:单独确认,批量确认以及异步确认;
三种Confirm实现模式对比:
| 模式 | 实现方式 | 性能 | 特点 |
|---|---|---|---|
| 单独确认 | 每发一条,调用waitForConfirms()等待确认 | 差 | 同步等待,吞吐量低 |
| 批量确认 | 发送一批消息,统一调用waitForConfirms() | 中 | 批量确认,提升吞吐量 |
| 异步确认 | 注册监听回调,异步处理确认结果 | 优 | 非阻塞,性能最高 |
异步确认代码实现(RabbitMQ):
/**
* 异步确认
* @throws IOException
* @throws TimeoutException
* @throws InterruptedException
*/
private static void handlingPublisherConfirmsAsynchronously() throws IOException, TimeoutException, InterruptedException {
try (Connection connection = getConnection()){
Channel channel = connection.createChannel();
channel.queueDeclare(Constants.PUBLISHER_CONFIRMS_QUEUE3, true, false, false, null);
// 1. 开启发布确认
channel.confirmSelect();
// 2. 使用一个集合, 保存未确认消息的 deliveryTag
// 集合必须有序, 因为需要批量确认
// 集合中存放的是未确认的 deliveryTag
SortedSet<Long> confirmSet = Collections.synchronizedSortedSet(new TreeSet<>());
// 3. 添加异步监听
channel.addConfirmListener(new ConfirmListener() {
@Override
public void handleAck(long deliveryTag, boolean multiple) throws IOException {
// 3.1 如果设置了批量确认
if (multiple){
// 删除 deliveryTag 之前的, 包含 deliveryTag, 表示已经确认 deliveryTag 及其之前的消息
confirmSet.headSet(deliveryTag + 1).clear();
}else { // 3.2 未设置批量确认
// 只删除 deliveryTag
confirmSet.remove(deliveryTag);
}
}
@Override
public void handleNack(long deliveryTag, boolean multiple) throws IOException {
// 3.3 如果设置了批量确认
if (multiple){
// 删除 deliveryTag 之前的, 包含 deliveryTag
confirmSet.headSet(deliveryTag + 1).clear();
}else { // 3.4 未设置批量确认
// 只删除 deliveryTag
confirmSet.remove(deliveryTag);
}
// 3.4 如果得到未确认, 需要根据具体的业务逻辑, 进行重发, 或者其它处理方式
// ...
}
});
// 4. 发送消息
long start = System.currentTimeMillis();
String message = "hello, publisher confirms...";
for (int i = 0; i < messageCount; i++){
long nextPublishSeqNo = channel.getNextPublishSeqNo();
channel.basicPublish("", Constants.PUBLISHER_CONFIRMS_QUEUE1, null, (message + " - " + i).getBytes());
// 发送消息后, 要把序号加入到未确认集合中
confirmSet.add(nextPublishSeqNo);
}
// 5. 消息确认完毕后, 记录时间
while(!confirmSet.isEmpty()){
Thread.sleep(10);
}
long end = System.currentTimeMillis();
System.out.println("共计发送 " + messageCount + " 条消息, 耗时: " + (end - start) + "ms");
}
}
异步确认代码实现(Spring):
配置:
spring:
## mq ##
rabbitmq:
publisher-confirm-type: correlated #消息发送确认
import org.springframework.amqp.core.ReturnedMessage;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitTemplateConfig {
@Bean("confirmRabbitTemplate")
public RabbitTemplate confirmRabbitTemplate(ConnectionFactory connectionFactory){
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
// 消息确认的回调方法
rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
@Override
public void confirm(CorrelationData correlationData, boolean b, String s) {
System.out.println("执行了confirm方法");
if (b){
System.out.printf("接收到消息, 消息id: %s",
correlationData == null ? null : correlationData.getId());
}else {
System.out.printf("未接收到消息, 消息id: %s, 原因: %s",
correlationData == null ? null : correlationData.getId(), s);
// 未接收到消息, 执行例如重发等业务逻辑
}
System.out.println();
}
});
return rabbitTemplate;
}
}
这是消息从生产者到交换机的过程,交换机确认后,也并不能保证消息就不会丢失,因为交换机收到消息后还需要交给队列,如果没有成功交给队列,消息就会被丢弃,造成了消息丢失后果;
因此,交换机确认后,如果没有把消息成功路由到队列中,消息应该退回给生产者,生产者根据业务逻辑,进行重发或者其它操作,以此保证消息的可靠传输;
2. Return 退回模式
消息到达交换机后,会根据路由规则进行匹配,把消息放入队列中。从交换机到队列的过程中,如果消息无法被路由到队列,可以选择把消息退回给生产者;
消息退回实现:
import org.springframework.amqp.core.ReturnedMessage;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitTemplateConfig {
@Bean("confirmRabbitTemplate")
public RabbitTemplate confirmRabbitTemplate(ConnectionFactory connectionFactory){
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
// 消息退回的回调方法
rabbitTemplate.setMandatory(true);
rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() {
@Override
public void returnedMessage(ReturnedMessage returnedMessage) {
System.out.println("消息退回: " + returnedMessage);
}
});
return rabbitTemplate;
}
}
结合Confirm确认以及Return退回模式,生产者能够最大限度确保消息可靠传输,如果消息丢失,生产者能够及时发现,后续会采取其它策略,尽可能确保消息从生产者到RabbitMQ队列中;
三、持久化
1. 持久化的三要素
RabbitMQ的消息持久化需要同时满足三个条件:
1. 交换机持久化
// Java API
channel.exchangeDeclare("durable.exchange", "direct", true); // durable=true
// Spring AMQP
@Bean("durableExchange")
public DirectExchange durableExchange(){
return ExchangeBuilder.directExchange("durable.exchange").build();
}
交换机持久化确保RabbitMQ重启后,交换机元数据不会丢失。
2. 队列持久化
// Java API
channel.queueDeclare("durable.queue", true, false, false, null); // durable=true
// Spring AMQP
@Bean("durableQueue")
public Queue durableQueue() {
return QueueBuilder.durable("durable.queue").build();
}
队列持久化确保队列元数据在重启后依然存在。
3. 消息持久化
// Java API - 设置消息持久化
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.deliveryMode(2) // 2表示持久化,1表示非持久化
.build();
channel.basicPublish("exchange", "routingKey", properties, message.getBytes());
// Spring AMQP - 使用MessageProperties
rabbitTemplate.convertAndSend("exchange", "routingKey", message, message -> {
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
});
// 或者在发送时指定MessagePostProcessor
rabbitTemplate.convertAndSend("exchange", "routingKey", message, new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
}
});
消息持久化确保消息内容写入磁盘,重启后可以恢复。
重要:只有同时满足以上三个条件,消息才能实现真正的持久化。
2. 持久化存在的问题
将交换机,队列,消息都设置了持久化之后,能够保证消息在RabbitMQ服务器中不会丢失吗?
答案也是否定的,在消息存入RabbitMQ的缓存中后,还需一小段时间才能将消息存入磁盘。RabbitMQ不会为每条消息都进行同步存盘,而是会将消息先保存到操作系统缓存中,如果还没来得及落盘,就发生了宕机,重启等异常情况,那么这些消息就会丢失;
为了避免这个问题,可以通过集群的方式,主节点出现问题,可以切换到从节点,保证了服务器高可用,除非集群都出现问题。
四、消息确认
1. 消息确认机制
生产者发送消息之后,到达消费者,可能会有如下情况:
- a. 消息处理成功
- b. 消息处理异常
RabbitMQ向消费者发送消息之后,就会把这条消息删掉,那么情况b就会造成消息丢失。为了保证消息从队列可靠得到达消费者,RabbitMQ提供了消息确认机制;
消费者订阅队列时,可以指定 autoAck 参数,根据参数设置,消息确认机制分为以下两中情况:
自动确认:autoAck 为 true,RabbitMQ会把发送出去得消息设置为确认,然后从内存或者磁盘中删除,并不关注消费者是否成功被消费者处理,这种模式适用于对于消息可靠性要求不高得场景;
手动确认:autoAck 为 false,RabbitMQ会等待消费者显示调用 Basic.Ack 命令,回复确认后才从内存或者磁盘中删除消息,这种模式适用于对于消息可靠性要求高得场景;
如果 autoAck 设置为 false,对于RabbitMQ服务端,队列中得消息分为两个部分:
- a. 等待投递给消费者的消息;
- b. 消息已经投递给消费者,但消费者还没确认;
如果RabbitMQ一直没收到消费者的确认消息,并且消费者已经断开连接,则会安排消息重新进入队列,等待投递给下一个消费者,当然也有可能还是原来的消费者;
手动确认关键API说明:
| 方法 | 说明 | 使用场景 |
|---|---|---|
| basicAck(tag, multiple) | 确认消息 | 处理成功时调用 |
| basicNack(tag, multiple, requeue) | 否定确认,可批量 | 处理失败时调用 |
| basicReject(tag, requeue) | 拒绝消息,单条 | 处理失败时调用 |
Spring-AMQP 对消息确认机制提供了三种策略:
- AcknowledgeMode.NONE:这种模式下,消息⼀旦投递给消费者,不管消费者是否成功处理了消息,RabbitMQ就会自动确认消息,从RabbitMQ队列中移除消息,如果消费者处理消息失败, 消息可能会丢失;
- AcknowledgeMode.AUTO(默认):这种模式下,消费者在消息处理成功时会自动确认消息,但如果处理过程中抛出了异常,则不会确认消息;
- AcknowledgeMode.MANUAL:手动确认模式下, 消费者必须在成功处理消息后显式调用basicAck 方法来确认消息. 如果消息未被确认,RabbitMQ会认为消息尚未被成功处理, 并且会在消费者可用时重新投递该消息, 这 种模式提高了消息处理的可靠性,因为即使消费者处理消息后失败,消息也不会丢失, 而是可以被重新处理;
注意:务必使用try-catch包裹业务逻辑
确保任何异常都能捕获并正确处理,避免消息既未确认也未拒绝,导致消息积压在未确认状态。
配置:
spring:
## mq ##
rabbitmq:
listener:
simple:
acknowledge-mode: manual
消费者代码实现:
import com.example.demo.constant.Constants;
import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import java.io.IOException;
//@Component
public class Consumer {
@RabbitListener(queues = Constants.ACK_QUEUE)
public void handleEvent(Message message, Channel channel) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try{
System.out.println("接收到消息: " + new String(message.getBody(), "UTF-8") +
", deliveryTag = " + deliveryTag);
// 肯定确认
channel.basicAck(deliveryTag, false);
} catch(Exception e){
// 否定确认
// channel.basicReject(deliveryTag, true);
channel.basicNack(deliveryTag, false, true);
}
}
}
总结
消息可靠性不是一个单一的技术点,而是一个系统工程,需要从生产者、服务端、消费者三个维度全链路考虑。生产者通过发布确认模式,确保消息能够从生产者到达RabbitMQ服务器,RabbitMQ通过持久化的方式,尽可能保证消息不在服务器丢失,消费者通过消息确认机制保证消息能够被消费者正确消费。
更多推荐



所有评论(0)