互联网全景消息(2)之RabbitMq高阶使用

2024-08-31 18:36

本文主要是介绍互联网全景消息(2)之RabbitMq高阶使用,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

一、RabbitMQ消息可靠性保障

        消息的可靠性投递是使用消息中间件不可避免的问题,不管是Kafka、rocketMQ或者是rabbitMQ,那么在RabbitMQ中如何保障消息的可靠性呢?

        首先来看一下rabbitMQ的 架构图:

        首先从图里我们可以看到,消息投递的保障性主要从三个方面来解决:

  • 生产者;
  • Broker;
  • 消费者; 

1.1 生产者保障 

        生产者发送消息到broker时,要保障消息的可靠性,主要方案有以下两种:

  1. 生产者确认;
  2. 失败通知; 

         首先RabbitMQ生产者通过制定一个Exchange和routingkey把消息送达到某个队列中,然后消费者监听队列进行消费处理。但是在某些情况下,如果我们在发送消息,当前的exchange不存在或者指定的routingkey找不到对应的队列,这个时候如果要监听这种不可达的消息,就需要失败通知了。

1.1.1 交换器、队列、路由健的关系

        队列通过路由健(routingkey,某种规则)绑定到交换器中,生产者将消息发布到交换器中,交换器根据绑定的routingkey将消息路由到指定队列中,然后由订阅这个队列的消费者进行监听消费。

 

        此时就会存在一个问题,消息路由到了不存在的队列怎么办?一般情况下RabbitMQ会直接忽略,当这个消息不存在,也就是消息丢弃了。

        所以在不做任何配置的情况下,生产者是不知道消息是否真正达到rabbitMQ,也就是说消息发布不会返回任何消息给生产者。

1.1.2 失败通知 

        那如何保证我们消息发布的可靠性,这里我们就可以启动失败通知,在原生的编码中可以在发送消息的时候设置Mandatory,即可开启故障检测模式。

        注意:他只会让RabbitMQ向你通知失败,而不会通知成功,如果消息正确的路由到队列,则发布者不会收到任何通知。带来的问题就是无法确保发布消息一定是成功的,因为通知失败的消息可能会丢失。

1.1.2.1 实现方式

        spring配置方式:

spring:

        rabbitmq:

                # 消息在未被队列收到的情况下返回

                publisher-returns: true

         关键代码,注意需要发送者实现 ReturnCallback 接口方可实现失败通知

 

1.1.2.2 存在的问题 

        如果消息正确路由到队列,则发布者不会收到任何通知。带来的问题就是无法确保消息一定是成功的,因为通知失败的消息可能会丢失。

        这样子我们可以使用RabbitMQ的发送方确认来实现,它不仅仅在路由失败的时候给我们发送消息,并且能够在路由成功的时候也给我们发送消息。

 1.1.3 发送方确认

        发送方确认是指生产者在投递消息后,如果Broker接收到消息,则会给生产者一个应答。生产者进行接收应答,用来确认这条消息是否正常的发送到Broker,这种方式也是可靠消息投递的核心保证。 

        rabbitMQ消息发送分为两个阶段:

  • 将消息发送到broker,即发送到exchange交换机;
  • 消息通过exchange被路由到队列; 

        一旦消息投递到队列,队列则会向生产者发送一个通知,如果队列设置了消息持久到磁盘,则会等待消息持久化到磁盘之后再发送通知。

        注意:发送者确认只有出现RabbitMQ内部错误才会出现发送者确认失败。 

        在发送者确认这种模式也可以分为具体两种情况来看待:

  1. 队列不可路由;
  2. 队列可路由; 
1.1.3.1 队列不可路由 

        当前的消息达到交换器之后,对于发送者确认是成功的。因为此时的消息已经到达了broker,此时只是不可路由队列他认为是成功的。

 

        首先RabbitMQ交换器不可路由时,消息也根本 不会投递到队列中,所以这里他只管到交换器这里,当消息成功到达交换器后,就会进行确认操作。 

        另外在这过程中,生产者收到了确认之后,那么因为消息不可路由,所以该消息也是无效的相当于被抛弃了,无法到达队列,所以一般这里会结合失败通知来一同使用,这里一般会进行mandatory模式,失败则会调用addReturnListener监听器来处理。

1.1.3.2 队列可以路由

        只要消息能够到达队列即可进行确认,一般是RabbitMQ发生内部错误才会出现确认失败的情况; 

         

        可以路由的消息,要等到被投递到所有匹配的队列之后,broker会发送一个确认给生产者(包含消息的唯一ID),这就使得生产者知道消息已经正确到达队列了。

        如果消息和队列是可持久化的,那么消息会在将消息写入磁盘之后发出,broker回传给生产者的确认小学中delivery-tag包含了确认消息的序列号。

1.1.3.3 使用方式

        Spring配置:

spring:rabbitmq:publisher-confirm-type: correlated

         关键代码,注意需要发送者实现 ConfirmCallback 接口方可实现失败通知:

 1.1.4 broker丢失消息

        前面我们从生产者的角度分析了消息可靠性传输的原理和实现,接下来就要看下broker是如何保障消息的可靠性传输的。

        假设生产者已经成功将消息发送到了交换机,并且交换机也成功的将消息路由到队列中,但是此时消费者还没有进行消费的时候,mq挂掉了,那么重启之后消息就会不存在,那样子就不能保障消息的可靠性 传输了。

        所以此时就要开启RabbitMQ的持久化,也就是将消息持久化到磁盘,此时即使MQ挂掉了,重启之后也会自动读取之前存储的数据。

1.1.4.1 持久化队列 

         在spring开启一个持久化队列。

  @Configurationpublic class RabbitConfig {public static final String DURABLE_QUEUE_NAME = "durable_queue";@Beanpublic Queue durableQueue() {// 创建一个持久化的队列return new Queue(DURABLE_QUEUE_NAME, true); // 第二个参数为true表示队列持久化}}
1.1.4.2 持久化交换器

@Configuration
public class RabbitConfig {public static final String DURABLE_EXCHANGE_NAME = "durable_exchange";public static final String DURABLE_QUEUE_NAME = "durable_queue";public static final String ROUTING_KEY = "durable_routing_key";@Beanpublic DirectExchange durableExchange() {// 创建一个持久化的Direct Exchangereturn new DirectExchange(DURABLE_EXCHANGE_NAME, true, false);}@Beanpublic Queue durableQueue() {// 创建一个持久化的队列return new Queue(DURABLE_QUEUE_NAME, true); // 第二个参数为true表示队列持久化}@Beanpublic Binding binding(Queue durableQueue, DirectExchange durableExchange) {// 绑定队列到交换器return BindingBuilder.bind(durableQueue).to(durableExchange).with(ROUTING_KEY);}
}
 1.1.4.3 发送持久化消息

         在发送消息的时候,需要设置属性deliveryMode=2,表示发送的是一个持久化消息,需要注意的是在springboot中,发送消息时已经自动设置了deliveryMode为2,不需要人工再去设置一遍。

@Component
public class MessageSender {@Autowiredprivate RabbitTemplate rabbitTemplate;public void sendPersistentMessage(String messageContent) {// 创建消息属性,并设置为持久化MessageProperties messageProperties = new MessageProperties();messageProperties.setDeliveryMode(MessageDeliveryMode.PERSISTENT);// 创建消息Message message = new Message(messageContent.getBytes(), messageProperties);// 发送消息到指定的交换器rabbitTemplate.convertAndSend(RabbitConfig.DURABLE_EXCHANGE_NAME, RabbitConfig.ROUTING_KEY, message);System.out.println("Sent message: " + messageContent);}
}

 1.1.5 总结

        生产者以及Broker要保障消息传递的可靠性如果结合失败通知以及发送方确认和持久化消息来实现。

1.发送方确认:保障消息能够到达broker;

2.失败通知:保障的是消息能够成功路由到队列;

3.持久化队列:保障消息的持久化;

1.2 消费者消息可靠性 

        消费者接收到消息,但是还未处理或者还未处理完成,此时消费者进程挂了,比如重启或者异常中断,此时mq会认为消费者已经完成消息消费,就会从队列中删除消息,从而导致消息丢失。 

        那该如何避免这种情况呢?这就要使用到RabbitMQ提供的ack机制,RabbitMQ默认是自动ack的,此时需要将其修改为手动ack,也就是自己的程序确定消息是否已经处理完成。如果此时出现消息未处理完成进程挂掉的情况,由于没有提交ack,rabbitMQ就不会删除这条消息,而是会把这条消息发送给其他消费者处理,但是消息不回丢失。 

spring:rabbitmq:listener:simple:acknowledge-mode: manual

        acknowledge-mode: manual代表开启手动ack,该配置项的其他两个参数值为none和auto;

  • auto:消费者根据程序执行的正常或者抛出异常来决定是抛出ack或者nack;
  • munual:手动ack,用户必须手动提交ack或者nack;
  • none:没有ack机制; 

        默认值是none,如果将ack的模式设置auto,此时如果消费者执行异常的话,就相当于执行了nack方法,消息会被放置到队列的头部,消息会被无限期的执行,从而导致后续消息无法执行。

这篇关于互联网全景消息(2)之RabbitMq高阶使用的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1124677

相关文章

JAVA系统中Spring Boot应用程序的配置文件application.yml使用详解

《JAVA系统中SpringBoot应用程序的配置文件application.yml使用详解》:本文主要介绍JAVA系统中SpringBoot应用程序的配置文件application.yml的... 目录文件路径文件内容解释1. Server 配置2. Spring 配置3. Logging 配置4. Ma

Linux使用dd命令来复制和转换数据的操作方法

《Linux使用dd命令来复制和转换数据的操作方法》Linux中的dd命令是一个功能强大的数据复制和转换实用程序,它以较低级别运行,通常用于创建可启动的USB驱动器、克隆磁盘和生成随机数据等任务,本文... 目录简介功能和能力语法常用选项示例用法基础用法创建可启动www.chinasem.cn的 USB 驱动

C#使用yield关键字实现提升迭代性能与效率

《C#使用yield关键字实现提升迭代性能与效率》yield关键字在C#中简化了数据迭代的方式,实现了按需生成数据,自动维护迭代状态,本文主要来聊聊如何使用yield关键字实现提升迭代性能与效率,感兴... 目录前言传统迭代和yield迭代方式对比yield延迟加载按需获取数据yield break显式示迭

使用SQL语言查询多个Excel表格的操作方法

《使用SQL语言查询多个Excel表格的操作方法》本文介绍了如何使用SQL语言查询多个Excel表格,通过将所有Excel表格放入一个.xlsx文件中,并使用pandas和pandasql库进行读取和... 目录如何用SQL语言查询多个Excel表格如何使用sql查询excel内容1. 简介2. 实现思路3

java脚本使用不同版本jdk的说明介绍

《java脚本使用不同版本jdk的说明介绍》本文介绍了在Java中执行JavaScript脚本的几种方式,包括使用ScriptEngine、Nashorn和GraalVM,ScriptEngine适用... 目录Java脚本使用不同版本jdk的说明1.使用ScriptEngine执行javascript2.

c# checked和unchecked关键字的使用

《c#checked和unchecked关键字的使用》C#中的checked关键字用于启用整数运算的溢出检查,可以捕获并抛出System.OverflowException异常,而unchecked... 目录在 C# 中,checked 关键字用于启用整数运算的溢出检查。默认情况下,C# 的整数运算不会自

在MyBatis的XML映射文件中<trim>元素所有场景下的完整使用示例代码

《在MyBatis的XML映射文件中<trim>元素所有场景下的完整使用示例代码》在MyBatis的XML映射文件中,trim元素用于动态添加SQL语句的一部分,处理前缀、后缀及多余的逗号或连接符,示... 在MyBATis的XML映射文件中,<trim>元素用于动态地添加SQL语句的一部分,例如SET或W

Mybatis官方生成器的使用方式

《Mybatis官方生成器的使用方式》本文详细介绍了MyBatisGenerator(MBG)的使用方法,通过实际代码示例展示了如何配置Maven插件来自动化生成MyBatis项目所需的实体类、Map... 目录1. MyBATis Generator 简介2. MyBatis Generator 的功能3

Python中使用defaultdict和Counter的方法

《Python中使用defaultdict和Counter的方法》本文深入探讨了Python中的两个强大工具——defaultdict和Counter,并详细介绍了它们的工作原理、应用场景以及在实际编... 目录引言defaultdict的深入应用什么是defaultdictdefaultdict的工作原理

使用Python进行文件读写操作的基本方法

《使用Python进行文件读写操作的基本方法》今天的内容来介绍Python中进行文件读写操作的方法,这在学习Python时是必不可少的技术点,希望可以帮助到正在学习python的小伙伴,以下是Pyth... 目录一、文件读取:二、文件写入:三、文件追加:四、文件读写的二进制模式:五、使用 json 模块读写