消息可靠性保证

回顾RabbitMQ的消息传递过程

image.png
如图所示,发生消息丢失的可能阶段也就是生产者发送消息,时rabbitmq存储消息时,消费者消费消息时。
项目源码:gitee

生产者发送消息阶段

  1. 生产者发送消息时把交换机名写错
  2. 生产者发送消息时把routingKey写错

RabbitMQ存储消息阶段

默认情况下rabbitmq会把消息存储到内存中,如果在消费者消费消息之前,rabbitmq服务器宕机了,内存就会被释放,消息就会丢失

消费者消息消息阶段

消费者在获取到消息以后,就会自动给rabbitmq服务端返回一个ack标志,rabbitmq服务端就会把这个消息从队列中删除。但当消费者获取到消息以后,准备进行业务逻辑处理时消费者宕机了,相当于该消息没有被消费成功,即消息丢失。

因此,我们就针对以上3个阶段,分别解决

生产者保证消息不丢失

  1. 生产者确认机制:可以让生产者感知到消息是否正常发送给交换机
  2. 生产者回退机制:可以让生产者感知到消息是否正常发送给队列

生产者确认机制

image.png

  1. 首先准备好环境,交换机,队列,绑定信息。
package com.example.rabbitmqreliable.demos;import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;@Configuration
public class RabbitMQConfig {public static final String CONFIRM_EXCHANGE_NAME = "confirm_exchange";public static final String CONFIRM_QUEUE_NAME = "confirm_queue";public static final String CONFIRM_ROUTING_KEY = "key1";@Beanpublic DirectExchange directExchange() {return new DirectExchange(CONFIRM_EXCHANGE_NAME);}@Beanpublic Queue confirmQueue() {return QueueBuilder.durable(CONFIRM_QUEUE_NAME).build();}@Beanpublic Binding confirmBind(@Qualifier("directExchange") DirectExchange confirmExchange,@Qualifier("confirmQueue") Queue confirmQueue) {return BindingBuilder.bind(confirmQueue).to(confirmExchange).with(CONFIRM_ROUTING_KEY);}
}
  1. 配置文件
spring.rabbitmq.host=101.133.141.75
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.rabbitmq.virtual-host=ConFirm
# 开启生产者确认机制,当消费者成功处理这个消息时,会向生产者发送一个确认信号,
# 告诉生产者这个消息已经被成功消费了。
# 如果生产者在一定时间内没有收到确认信号,就会重新发送这个消息。
spring.rabbitmq.publisher-confirm-type=correlated
  1. 通过测试类,创建生产者发送消息
package com.example.rabbitmqreliable;import com.example.rabbitmqreliable.demos.RabbitMQConfig;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;@SpringBootTest
public class ProviderTests {/*** 目标:让生产者获取到rabbitmq服务返回的ack或nack* 做法:rabbitTemplate需要绑定对应的回调函数* 分析:目前的rabbitTemplate是spring托管的,并没有对应的回调函数,需要自定义* 实施:需要自定义rabbitTemplate,并注入到spring容器中* 一旦我们在spring容器中配置了一个rabbitTemplate,* 那么spring boot就不会对rabbitTemplate进行自动化配置*/@Autowiredprivate RabbitTemplate rabbitTemplate;public void test1() {rabbitTemplate.convertAndSend(RabbitMQConfig.CONFIRM_EXCHANGE_NAME, RabbitMQConfig.CONFIRM_ROUTING_KEY, "HELLO CONFIRM");}}
  1. 自定义rabbitTemplate,实现确认机制的回调方法,需要在RabbitMQConfig文件中添加以下内容:
/*** ConnectFactory由spring boot根据配置文件中的连接信息实现自动化配置* 即在spring容器中直接存在了ConnectionFactory对象* @param connectionFactory* @return*/@Beanpublic RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);// 设置回调函数// 而ConfirmCallBack是一个接口,需要一个类去实现他rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {/*** 当rabbitmq服务端给生产者放回ack/nack时会执行该方法* @param correlationData 消息的id,内容* @param ack 消息是否发送成功* @param cause 原因*/@Overridepublic void confirm(CorrelationData correlationData, boolean ack, String cause) {if(ack) {System.out.println("消息正常发送给交换机");}else {System.out.println("消息没有正常发送给交换机,cause" + cause);// TODO 处理方案:再次发送消息给rabbitmq,需要获取消息内容}}});return rabbitTemplate;}
  1. 实现当消息发送失败时,再次重新发送部分。
rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {/*** 当rabbitmq服务端给生产者放回ack/nack时会执行该方法* @param correlationData 消息的id,内容* @param ack 消息是否发送成功* @param cause 原因*/@Overridepublic void confirm(CorrelationData correlationData, boolean ack, String cause) {if(ack) {System.out.println("消息正常发送给交换机");}else {System.out.println("消息没有正常发送给交换机,cause" + cause);// 处理方案:再次发送消息给rabbitmq,需要获取消息内容//方案1: 立马拿着id去数据库查消息// 方案2:通过定时任务重新发送String msgId = correlationData.getId();System.out.println("msgId" + msgId);// 规定消息的最大发送次数3次,发送消息前判断实际发送次数是否大于最大发送次数,如果大于就不进行重新发送,并设置status=2}}});
 @Testpublic void test1() {// 发送消息前把消息写入数据库,并分配唯一id(如果发送失败,可以拿这个id去查数据库,重新发送)// 并且还要记录消息实际发送次数,以及消息状态。当超过发送次数超过了规定值,就设置消息的status为2,发送成功设置消息状态为1String msgId = UUID.randomUUID().toString().replace("-","");CorrelationData correlationData =new CorrelationData(msgId);rabbitTemplate.convertAndSend(RabbitMQConfig.CONFIRM_EXCHANGE_NAME, RabbitMQConfig.CONFIRM_ROUTING_KEY + "error", "HELLO CONFIRM", correlationData);}

image.png

生产者回退机制

image.png

  1. 需要在配置文件中开启生产者回退机制
# 开启生产者确认机制,
spring.rabbitmq.publisher-returns=true
  1. 给rabbitTemplate绑定生产者回退机制的回调函数
/*** 给rabbitTemplate绑定回退机制的回调函数* ReturnCallback是一个接口,使用匿名内部类实现* 该方法被调用的概率极低,因为从交换机到队列的过程是rabbitmq内部实现的* 如果会出错,咱们也不会用他*/rabbitTemplate.setMandatory(true);//让rabbitmq服务把失败信息回传给生产者rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() {// 当消息没有正常转发给队列的时候被调用@Overridepublic void returnedMessage(ReturnedMessage returnedMessage) {byte[] body = returnedMessage.getMessage().getBody();String msg = new String(body);System.out.println("msg:" + msg);}});
  1. 执行测试方法
package com.example.rabbitmqreliable;import com.example.rabbitmqreliable.demos.RabbitMQConfig;
import org.junit.jupiter.api.Test;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;import java.util.UUID;@SpringBootTest
public class ProviderTests {/*** 目标:让生产者获取到rabbitmq服务返回的ack或nack* 做法:rabbitTemplate需要绑定对应的回调函数* 分析:目前的rabbitTemplate是spring托管的,并没有对应的回调函数,需要自定义* 实施:需要自定义rabbitTemplate,并注入到spring容器中* 一旦我们在spring容器中配置了一个rabbitTemplate,* 那么spring boot就不会对rabbitTemplate进行自动化配置*/@Autowiredprivate RabbitTemplate rabbitTemplate;@Testpublic void test1() {// 发送消息前把消息写入数据库,并分配唯一id(如果发送失败,可以拿这个id去查数据库,重新发送)// 并且还要记录消息实际发送次数,以及消息状态。当超过发送次数超过了规定值,就设置消息的status为2,发送成功设置消息状态为1String msgId = UUID.randomUUID().toString().replace("-","");CorrelationData correlationData =new CorrelationData(msgId);// 把routingKey写错rabbitTemplate.convertAndSend(RabbitMQConfig.CONFIRM_EXCHANGE_NAME, RabbitMQConfig.CONFIRM_ROUTING_KEY + "404", "HELLO CONFIRM", correlationData);}}

image.png

RabbitMQ保证消息不丢失

  1. 对交换机进行持久化
  2. 对队列进行持久化
  3. 对消息进行持久化

消费者保证消息不丢失

Spring Boot整合RabbitMQ消费者的应答模式:

  • none:自动应答,消费者获取到消息后直接给rabbitmq返回ack
  • auto(默认值):由spring boot框架根据业务执行特点决定给rabbitmq返回ack还nack,业务正常执行完毕返回ack,业务执行过程产生异常,返回nack
  • manual:手动应答,由程序员自己根据业务执行特点给rabbitmq返回对应的ack或nack

使用none模式

  1. 在配置文件中设置消费者应答模式
# 消费者的应答模式
spring.rabbitmq.listener.simple.acknowledge-mode=none
  1. 设置消费者
@Component
public class Consumer {@RabbitListener(queues = RabbitMQConfig.CONFIRM_QUEUE_NAME) // 监听的队列public void consumerListener(Message message) {byte[] body = message.getBody();;String msg = new String(body);// 进行业务处理int a = 1 / 0; //产生异常System.out.println("[consumerListener],msg:" + msg);}
}
  1. 启动项目,再运行测试类

image.png
image.png
由于我们使用none自动应答模式,消费者给rabbitmq返回ack,rabbitmq直接把消息从队列中删除,导致消息丢失

使用auto模式

修改配置文件中的应答方式为auto,以及引伸出来的其他配置项

# 消费者的应答模式
spring.rabbitmq.listener.simple.acknowledge-mode=auto
# 开启重试机制
spring.rabbitmq.listener.simple.retry.enabled=true
# 最大重试次数,否则会无限重试下去
spring.rabbitmq.listener.simple.retry.max-attempts=3
# 初始化的重试时间间隔
spring.rabbitmq.listener.simple.retry.initial-interval=1000
# 最大的重试时间间隔
spring.rabbitmq.listener.simple.retry.max-interval=5000
# 乘子(计算每一次时间间隔):1s->2s->4s->5s
spring.rabbitmq.listener.simple.retry.multiplier=2

此时控制台报了3次错误,队列没有消息,消息还是丢失了。

auto模式需要设置最大重试次数,否则会死循环,但是又无法判断最大重试次数是多少

使用manual模式

  1. 修改配置文件中的应答方式为manual,将auto模式中延申出来的配置项注释掉
spring.rabbitmq.listener.simple.acknowledge-mode=manual
  1. 消费者代码
package com.example.rabbitmqreliable.demos;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 = RabbitMQConfig.CONFIRM_QUEUE_NAME) // 监听的队列public void consumerListener(Message message, Channel channel) {byte[] body = message.getBody();;String msg = new String(body);long deliveryTag = message.getMessageProperties().getDeliveryTag();try {// 进行业务处理int a = 1 / 0; //产生异常System.out.println("[consumerListener],msg:" + msg);// 没有产生异常,给服务端返回ack// 第一个参数:表示消息的标签,保证消息唯一性// 第二个数:表示是否需要进行批量应答channel.basicAck(deliveryTag, true);} catch (Exception e) {e.printStackTrace();// 产生异常// 给rabbitmq返回nack// 第三个参数:表示是否将消息重新放入队列中try {channel.basicNack(deliveryTag, true, true);} catch (IOException ex) {throw new RuntimeException(ex);}}}
}

启动项目后执行测试代码,控制台会不断报红,死循环。因为没有设置最大重试次数,因此我们需要统计消息的实际消费次数,可以借助redis计算。一旦消息的实际消费次数大于最大消费次数,那么此时需要给rabbitmq返回ack删除该消息,返回之前要将该消息记录数据库中,后期人工处理