目录

前言

一、消息丢失的三个环节

1. 消息的生命周期

2. 丢失消息三个环节

3. 消息可靠投递的设计思路

二、发布确认模式

1. Confirm确认模式

2. Return 退回模式

三、持久化

1. 持久化的三要素

2. 持久化存在的问题

四、消息确认

1. 消息确认机制

总结


前言

在分布式系统中,消息队列作为核心组件,承担着系统解耦、流量削峰的重任。而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通过持久化的方式,尽可能保证消息不在服务器丢失,消费者通过消息确认机制保证消息能够被消费者正确消费。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐