互联网全景消息(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

相关文章

中文分词jieba库的使用与实景应用(一)

知识星球:https://articles.zsxq.com/id_fxvgc803qmr2.html 目录 一.定义: 精确模式(默认模式): 全模式: 搜索引擎模式: paddle 模式(基于深度学习的分词模式): 二 自定义词典 三.文本解析   调整词出现的频率 四. 关键词提取 A. 基于TF-IDF算法的关键词提取 B. 基于TextRank算法的关键词提取

使用SecondaryNameNode恢复NameNode的数据

1)需求: NameNode进程挂了并且存储的数据也丢失了,如何恢复NameNode 此种方式恢复的数据可能存在小部分数据的丢失。 2)故障模拟 (1)kill -9 NameNode进程 [lytfly@hadoop102 current]$ kill -9 19886 (2)删除NameNode存储的数据(/opt/module/hadoop-3.1.4/data/tmp/dfs/na

Hadoop数据压缩使用介绍

一、压缩原则 (1)运算密集型的Job,少用压缩 (2)IO密集型的Job,多用压缩 二、压缩算法比较 三、压缩位置选择 四、压缩参数配置 1)为了支持多种压缩/解压缩算法,Hadoop引入了编码/解码器 2)要在Hadoop中启用压缩,可以配置如下参数

Makefile简明使用教程

文章目录 规则makefile文件的基本语法:加在命令前的特殊符号:.PHONY伪目标: Makefilev1 直观写法v2 加上中间过程v3 伪目标v4 变量 make 选项-f-n-C Make 是一种流行的构建工具,常用于将源代码转换成可执行文件或者其他形式的输出文件(如库文件、文档等)。Make 可以自动化地执行编译、链接等一系列操作。 规则 makefile文件

使用opencv优化图片(画面变清晰)

文章目录 需求影响照片清晰度的因素 实现降噪测试代码 锐化空间锐化Unsharp Masking频率域锐化对比测试 对比度增强常用算法对比测试 需求 对图像进行优化,使其看起来更清晰,同时保持尺寸不变,通常涉及到图像处理技术如锐化、降噪、对比度增强等 影响照片清晰度的因素 影响照片清晰度的因素有很多,主要可以从以下几个方面来分析 1. 拍摄设备 相机传感器:相机传

【专题】2024飞行汽车技术全景报告合集PDF分享(附原数据表)

原文链接: https://tecdat.cn/?p=37628 6月16日,小鹏汇天旅航者X2在北京大兴国际机场临空经济区完成首飞,这也是小鹏汇天的产品在京津冀地区进行的首次飞行。小鹏汇天方面还表示,公司准备量产,并计划今年四季度开启预售小鹏汇天分体式飞行汽车,探索分体式飞行汽车城际通勤。阅读原文,获取专题报告合集全文,解锁文末271份飞行汽车相关行业研究报告。 据悉,业内人士对飞行汽车行业

pdfmake生成pdf的使用

实际项目中有时会有根据填写的表单数据或者其他格式的数据,将数据自动填充到pdf文件中根据固定模板生成pdf文件的需求 文章目录 利用pdfmake生成pdf文件1.下载安装pdfmake第三方包2.封装生成pdf文件的共用配置3.生成pdf文件的文件模板内容4.调用方法生成pdf 利用pdfmake生成pdf文件 1.下载安装pdfmake第三方包 npm i pdfma

零基础学习Redis(10) -- zset类型命令使用

zset是有序集合,内部除了存储元素外,还会存储一个score,存储在zset中的元素会按照score的大小升序排列,不同元素的score可以重复,score相同的元素会按照元素的字典序排列。 1. zset常用命令 1.1 zadd  zadd key [NX | XX] [GT | LT]   [CH] [INCR] score member [score member ...]

git使用的说明总结

Git使用说明 下载安装(下载地址) macOS: Git - Downloading macOS Windows: Git - Downloading Windows Linux/Unix: Git (git-scm.com) 创建新仓库 本地创建新仓库:创建新文件夹,进入文件夹目录,执行指令 git init ,用以创建新的git 克隆仓库 执行指令用以创建一个本地仓库的

【北交大信息所AI-Max2】使用方法

BJTU信息所集群AI_MAX2使用方法 使用的前提是预约到相应的算力卡,拥有登录权限的账号密码,一般为导师组共用一个。 有浏览器、ssh工具就可以。 1.新建集群Terminal 浏览器登陆10.126.62.75 (如果是1集群把75改成66) 交互式开发 执行器选Terminal 密码随便设一个(需记住) 工作空间:私有数据、全部文件 加速器选GeForce_RTX_2080_Ti