消息中间件——RabbitMQ(七)高级特性 2

前言

上一篇消息中间件——RabbitMQ(七)高级特性 1中我们介绍了消息如何保障100%的投递成功?,幂等性概念详解,在海量订单产生的业务高峰期,如何避免消息的重复消费的问题?,Confirm确认消息、Return返回消息。这篇我们来介绍下下面内容。

  • 自定义消费者
  • 消息的限流(防止占用内存过多,节点宕机)
  • 消息的ACK与重回队列
  • TTL消息
  • 死信队列

1. 自定义消费者

1.1 消费端自定义监听

我们一般就在代码中编写while循环,进行consumer.nextDelivery方法进行获取下一条消息,然后进行消费处理!

但是这种轮训的方式肯定是不好的,代码也比较low。

  • 我们使用自定义的Consumer更加的方便,解耦性更加的强,也是在实际工作中最常见的使用方式!

1.2 代码演示

1.2.1 生产者
public class Producer {public static void main(String[] args) throws Exception {//1 创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();String exchange = "test_consumer_exchange";String routingKey = "consumer.save";String msg = "Hello RabbitMQ Consumer Message";for(int i =0; i<5; i ++){channel.basicPublish(exchange, routingKey, true, null, msg.getBytes());}}
}
1.2.2 消费者
public class Consumer {public static void main(String[] args) throws Exception {// 创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();String exchangeName = "test_consumer_exchange";String routingKey = "consumer.#";String queueName = "test_consumer_queue";channel.exchangeDeclare(exchangeName, "topic", true, false, null);channel.queueDeclare(queueName, true, false, false, null);channel.queueBind(queueName, exchangeName, routingKey);//实现自己的MyConsumer()channel.basicConsume(queueName, true, new MyConsumer(channel));}
}
1.2.3 自定义类:MyConsumer
public class MyConsumer extends DefaultConsumer {public MyConsumer(Channel channel) {super(channel);}//根据需求,重写自己需要的方法。@Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.err.println("-----------consume message----------");//消费标签System.err.println("consumerTag: " + consumerTag);//这个对象包含许多关键信息System.err.println("envelope: " + envelope);System.err.println("properties: " + properties);System.err.println("body: " + new String(body));}}

1.3 打印结果

2. 消费端限流

2.1 什么是消费端的限流?

  • 假设一个场景,首先,我们Rabbitmq服务器有上万条未处理的消息,我们随便打开一个消费者客户端,会出现下面情况:
  • 巨量的消息瞬间全部推送过来,但是我们单个客户端无法同时处理这么多数据!这个时候很容易导致服务器崩溃,出现故障。

为什么不在生产端进行限流呢?

因为在高并发的情况下,客户量就是非常大,所以很难在生产端做限制。因此我们可以用MQ在消费端做限流。

  • RabbitMQ提供了一种qos(服务质量保证)功能,即在非自动确认消息的前提下,如果一定数目的消息(通过基于consume或者channel设置Qos的值)未被确认前,不进行消费新的消息。
    在限流的情况下,千万不要设置自动签收,要设置为手动签收
  • void BasicQos(uint prfetchSize,ushort prefetchCount,bool global);

参数解释:
prefetchSize:0
prefetchCount:会告诉RabbitMQ不要同时给一个消费者推送多于N个消息,即一旦有N个消息还没有ack,则该consumer将block掉,直到有消息ack。
global: true\false 是否将上面设置应用于channel,简单点说,就是上面限制是channel级别还是consumer级别。
prefetchSize和global这两项,rabbitmq没有实现,暂且不研究prefetch_count在no_ask = false的情况下生效,即在自动应答的情况下这两个值是不生效的。

第一个参数:消息的限制大小,消息多少兆。一般不做限制,设置为0
第二个参数:一次最多处理多少条,实际工作中设置为1就好
第三个参数:限流策略在什么上应用。在RabbitMQ一般有两个应用级别:1.通道 2.Consumer级别。一般设置为false,true 表示channel级别,false表示在consumer级别

2.2 代码演示

2.2.1 生产者
public class Producer {public static void main(String[] args) throws Exception {//1 创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();String exchange = "test_qos_exchange";String routingKey = "qos.save";String msg = "Hello RabbitMQ QOS Message";for(int i =0; i<5; i ++){channel.basicPublish(exchange, routingKey, true, null, msg.getBytes());}}
}
2.2.2 消费者
public class Consumer {public static void main(String[] args) throws Exception {		//1 创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();		String exchangeName = "test_qos_exchange";String queueName = "test_qos_queue";String routingKey = "qos.#";channel.exchangeDeclare(exchangeName, "topic", true, false, null);channel.queueDeclare(queueName, true, false, false, null);channel.queueBind(queueName, exchangeName, routingKey);//1 限流方式  第一件事就是 autoAck设置为 false//设置为1,表示一条一条数据处理channel.basicQos(0, 1, false);channel.basicConsume(queueName, false, new MyConsumer(channel));		}
}
2.2.3 自定义类:MyConsumer
public class MyConsumer extends DefaultConsumer {private Channel channel ;public MyConsumer(Channel channel) {super(channel);this.channel = channel;}@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.err.println("-----------consume message----------");System.err.println("consumerTag: " + consumerTag);System.err.println("envelope: " + envelope);System.err.println("properties: " + properties);System.err.println("body: " + new String(body));//需要做签收,false表示不支持批量签收channel.basicAck(envelope.getDeliveryTag(), false);}}
2.2.4 测试结果

我们先注释掉:channel.basicAck(envelope.getDeliveryTag(), false);然后启动Consumer。
查看Exchange

查看Queues

然后再启动Producer。查看打印结果:

我们会发现消费端,只收到了一条消息。这是为什么呢?

第一点因为我们在consumer中

channel.basicConsume(queueName, false, new MyConsumer(channel));

第二个参数设置为false为手动签收。

第二点在qos中设置只接受一条消息。如果这一条消息不给Broker Ack应答的话,那么Broker会认为你并没有消费完这一条消息,那么就不会继续发送消息。

channel.basicQos(0, 1, false);

可以看下管控台,unack=1,Ready=4,total=5.


接下来我们放开注释channel.basicAck(envelope.getDeliveryTag(), false); 进行消息签收。重启服务。

3.1 打印结果

可以看到正常打印五条结果

4. 消费端ACK与重回队列

4.1 消费端的手工ACK和NACK

消费端进行消费的时候,如果由于业务异常我们可以进行日志的记录,然后进行补偿!

如果由于服务器宕机等严重问题,那我们就需要手工进行ACK保障消费端消费成功!

4.2 消费端的重回队列

消费端重回队列是为了对没有处理成功的消息,把消息重新传递给Broker!

一般我们在实际应用中,都会关闭重回队列,也就是设置为False.

4.3 代码演示

4.3.1 生产者
public class Producer {public static void main(String[] args) throws Exception {//1创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();String exchange = "test_ack_exchange";String routingKey = "ack.save";		for(int i =0; i<5; i ++){Map<String, Object> headers = new HashMap<String, Object>();headers.put("num", i);//添加属性,后续会使用到AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder().deliveryMode(2) //投递模式,持久化.contentEncoding("UTF-8").headers(headers).build();String msg = "Hello RabbitMQ ACK Message " + i;channel.basicPublish(exchange, routingKey, true, properties, msg.getBytes());}}
}
4.3.2 消费者
public class Consumer {public static void main(String[] args) throws Exception {//1创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();String exchangeName = "test_ack_exchange";String queueName = "test_ack_queue";String routingKey = "ack.#";channel.exchangeDeclare(exchangeName, "topic", true, false, null);channel.queueDeclare(queueName, true, false, false, null);channel.queueBind(queueName, exchangeName, routingKey);// 手工签收 必须要关闭 autoAck = falsechannel.basicConsume(queueName, false, new MyConsumer(channel));}
}
4.3.3 自定义类:MyConsumer
public class MyConsumer extends DefaultConsumer {private Channel channel ;public MyConsumer(Channel channel) {super(channel);this.channel = channel;}@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.err.println("-----------consume message----------");System.err.println("body: " + new String(body));try {Thread.sleep(2000);} catch (InterruptedException e) {e.printStackTrace();}if((Integer)properties.getHeaders().get("num") == 0) {//Nack三个参数  第二个参数:是否是批量,第三个参数:是否重回队列(需要注意可能会发生重复消费,造成死循环)channel.basicNack(envelope.getDeliveryTag(), false, true);} else {channel.basicAck(envelope.getDeliveryTag(), false);}}}

5.1 打印结果:

注意:
可以看到重回队列会出现重复消费导致死循环的问题,这时候最好设置重试次数,比如超过三次后,消息还是消费失败,就将消息丢弃。

6. TTL队列/消息

6.1 TTL

  • TTL是Time To Live的缩写,也就是生存时间
  • RabbitMQ支持消息的过期时间,在消息发送时可以进行指定
  • RabbitMQ支持队列的过期时间,从消息入队列开始计算,只要超过了队列的超时时间配置,那么消息会自动的清除

6.2 代码演示

6.2.1 直接通过管控台进行演示

通过管控台创建一个队列

x-max-length 队列的最大大小
x-message-ttl 设置10秒钟,如果消息还没有被消费的话,就会被清除。

添加exchange

Queue与Exchange进行绑定

点击 test_ttl_exchange 进行绑定

查看是否绑定成功

通过管控台发送消息

消息未处理自动清除

生产端设置过期时间

AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder().deliveryMode(2).contentEncoding("UTF-8").expiration("10000").headers(headers).build();

这两个属性并不相同,一个对应的是消息体,一个对应的是队列的过期。

7. 死信队列

7.1 概念理解

死信队列:DLX,Dead-Letter-Exchange
RabbitMQ的死信队里与Exchange息息相关

  • 利用DLX,当消息在一个队列中变成死信(dead message)之后,它能被重新publish到另一个Exchange,这个Exchange就是DLX

消息变成死信有以下几种情况

  • 消息被拒绝(basic.reject/basic.nack)并且requeue=false
  • 消息TTL过期
  • 队列达到最大长度

DLX也是一个正常的Exchange,和一般的Exchange没有区别,它能在任何的队列上被指定,实际上就是设置某个队列的属性

当这个队列中有死信时,RabbitMQ就会自动的将这个消息重新发布到设置的Exchange上去,进而被路由到另一个队列。

可以监听这个队列中消息做相应的处理,这个特征可以弥补RabbitMQ3.0以前支持的immediate参数的功能。

7.2 代码演示

  • 死信队列设置:
  • 首先需要设置死信队列的exchange和queue,然后进行绑定:
    Exchange:dlx.exchange
    Queue:dlx.queue
    RoutingKey:#
  • 然后我们进行正常声明交换机、队列、绑定,只不过我们需要在队列加上一个参数即可:arguments.put(“x-dead-letter-exchange”,”dlx.exchange”);
  • 这样消息在过期、requeue、队列在达到最大长度时,消息就可以直接路由到死信队列!
7.2.1 生产者
public class Producer {public static void main(String[] args) throws Exception {//创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();String exchange = "test_dlx_exchange";String routingKey = "dlx.save";String msg = "Hello RabbitMQ DLX Message";for(int i =0; i<1; i ++){AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder().deliveryMode(2).contentEncoding("UTF-8").expiration("10000").build();channel.basicPublish(exchange, routingKey, true, properties, msg.getBytes());}}
}
7.2.2 消费者

public class Consumer {public static void main(String[] args) throws Exception {//创建ConnectionFactoryConnection connection = ConnectionUtils.getConnection();Channel channel = connection.createChannel();// 这就是一个普通的交换机 和 队列 以及路由String exchangeName = "test_dlx_exchange";String routingKey = "dlx.#";String queueName = "test_dlx_queue";channel.exchangeDeclare(exchangeName, "topic", true, false, null);Map<String, Object> agruments = new HashMap<String, Object>();agruments.put("x-dead-letter-exchange", "dlx.exchange");//这个agruments属性,要设置到声明队列上channel.queueDeclare(queueName, true, false, false, agruments);channel.queueBind(queueName, exchangeName, routingKey);//要进行死信队列的声明:channel.exchangeDeclare("dlx.exchange", "topic", true, false, null);channel.queueDeclare("dlx.queue", true, false, false, null);channel.queueBind("dlx.queue", "dlx.exchange", "#");channel.basicConsume(queueName, true, new MyConsumer(channel));}
}
7.2.3 自定义类:MyConsumer
public class MyConsumer extends DefaultConsumer {public MyConsumer(Channel channel) {super(channel);}@Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.err.println("-----------consume message----------");System.err.println("consumerTag: " + consumerTag);System.err.println("envelope: " + envelope);System.err.println("properties: " + properties);System.err.println("body: " + new String(body));}
}
7.2.4 测试结果

运行Consumer,查看管控台


查看Exchanges

查看queue

可以看到test_dlx_queue多了DLX的标识,表示当队列中出现死信的时候,会将消息发送到死信队列dlx_queue中

关闭Consumer,只运行Producer

过10秒钟后,消息过期

在我们工作中,死信队列非常重要,用于消息没有消费者,处于死信状态。我们可以才用补偿机制。

小结

本次主要介绍了RabbitMQ的高级特性,首先介绍了互联网大厂在实际使用中如何保障100%的消息投递成功和幂等性的,以及对RabbitMQ的确认消息、返回消息、ACK与重回队列、消息的限流,以及对超时时间、死信队列的使用

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.mzph.cn/news/182854.shtml

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

作为用户,推荐算法真的是最优解么?

前言 众所周知&#xff0c;随着互联网技术的发展&#xff0c;推荐算法也越来越普及。无论是购物网站、社交媒体平台还是在线影视平台&#xff0c;推荐算法已成为用户获取相关信息的主要途径。据悉&#xff0c;近期GitHub决定结合算法推荐&#xff0c;将“Following”和“For Yo…

uniapp打包ios有时间 uniapp打包次数

我们经常用的解决方案有,分包,将图片上传到服务器上,减少插件引入。但是还有一个方案好多刚入门uniapp的人都给忽略了,就是在源码视图中配置,开启分包优化。 1.分包 目前微信小程序可以分8个包,每个包的最大存储是2M,也就是说你文件总体的大小不能超过16M,每个包的大…

刷题学习记录

[SWPUCTF 2021 新生赛]sql 进入环境 查看源码&#xff0c;发现是get传参且参数为wllm fuzz测试&#xff0c;发现空格&#xff0c;&#xff0c;and被过滤了 同样的也可以用python脚本进行fuzz测试 import requests fuzz{length ,,handler,like,select,sleep,database,delete,h…

【MySQL】常用内置函数:数值函数 / 字符串函数 / 日期函数 / 其他函数

文章目录 数值函数round()&#xff1a;四舍五入ceiling()&#xff1a;上限函数floor()&#xff1a;地板函数abs()&#xff1a;计算绝对值rand()&#xff1a;生成0-1的随机浮点数 字符串函数length()&#xff1a;获取字符串中的字符数upper() / lower()&#xff1a;将字符串转化…

算法---定长子串中元音的最大数目

题目 给你字符串 s 和整数 k 。 请返回字符串 s 中长度为 k 的单个子字符串中可能包含的最大元音字母数。 英文中的 元音字母 为&#xff08;a, e, i, o, u&#xff09;。 示例 1&#xff1a; 输入&#xff1a;s “abciiidef”, k 3 输出&#xff1a;3 解释&#xff1a;…

ThinkPHP的方法接收json数据问题

第一次接触到前后端分离开发&#xff0c;需要在后端接收前端ajax提交的json数据&#xff0c;开发基于ThinkPHP3.2.3框架。于是一开始习惯性的直接用I()方法接收到前端发送的json数据&#xff0c;然后用json_decode()解析发现结果为空&#xff01;但是打印出还未解析的值却打印得…

【代码随想录算法训练营-第一天】【数组】704. 二分查找、27. 移除元素

LeetCode-704.二分查找 【错误】第一遍提交的代码 主要错误点&#xff1a; 没弄清楚区间的定义导致&#xff1a;r 在定义处的赋值和 if 判断之后 r 的复制没有想清楚&#xff1b;没有搞清楚判断循环结束的条件&#xff1b;没有搞明白区间的定义&#xff0c;r 和 l 如何赋值&a…

使用Ray轻松进行Python分布式计算

大家好&#xff0c;在实际研究中&#xff0c;即使是具有多个CPU核心的单处理器计算机&#xff0c;也会给人一种能够同时运行多个任务的错觉。当我们拥有多个处理器时&#xff0c;就可以真正以并行的方式执行计算&#xff0c;本文将简要介绍Python分布式计算。 1.并行计算与分布…

MySQL事务日志

文章目录 1. redo日志 1. redo日志 口述&#xff1a;redo log 日志其实保证了ACID中的持久性&#xff0c;就是说当事务commit后&#xff0c;那么相应的修改呀更新这些操作其实都会记录到redo log中&#xff0c;其实这里的操作还是区别于redis中的aof中&#xff0c;它不是具体的…

文件中找TopK问题

目录 1.解题思路2.创建一个文件并在文件中写入数据3.为什么要建立小堆而不建立大堆&#xff1f;4.如何在现有的数据中建立适合的大堆&#xff1f;5.代码实现 1.解题思路 TopK问题即是在众多数据中找出前K大的值&#xff0c;则可以根据堆的性质来实现&#xff0c;但在使用堆之前…

Guacamole简介及centos7下搭建教程

简介 Guacamole是一款开源的远程桌面框架&#xff0c;它允许用户通过Web浏览器远程访问计算机资源。 官网地址&#xff1a;Apache Guacamole™ 官方文档&#xff1a;Installing Guacamole natively — Apache Guacamole Manual v1.5.3 架构 组件描述客户端浏览器用户通过支…

数据结构 / day02作业

1. 有若干个学校人员的信息,包括学生和教师。 其中学生的数据包括&#xff1a;姓名、性别、职业s/S、分数。 教师的数据包括:姓名、性别、职业t/T、职务。 1&#xff0c;定义指针指向堆区内存 2.循环输入 3.计算老师的个数 4.计算学生的平均值 5.循环输出 6释放堆区空间 #inc…

Jquery中的bind(),live(),delegate(),on()绑定事件方式

bind() 简要描述 bind()向匹配元素添加一个或多个事件处理器。 使用方式 $(selector).bind(event,data,function) event&#xff1a;必需项&#xff1b;添加到元素的一个或多个事件&#xff0c;例如 click,dblclick等&#xff1b; 单事件处理&#xff1a;例如 $(selector).bi…

不同类型的开源许可证

不同类型的开源许可证 什么是开源许可证 最简单的解释是&#xff0c;开源许可证是计算机软件和其他产品的许可证&#xff0c;允许在定义的条款和条件下使用、修改或共享源代码、蓝图或设计。开源并不意味着该软件可以根据需要使用、复制、修改和分发。根据开源许可证的类型&a…

gcp, loki, honeybadger 查看日志 语句

gcp, loki, honeybadger 查看日志 语句 GCP resource.type“cloudsql_database” AND logName:“projects/ebay-mag/logs/cloudsql.googleapis.com%2Fpostgres.log” AND ( textPayload:“[2823409]”) log_name“projects/ebay-mag/logs/cloudsql.googleapis.com%2Fpostgres…

代码随想录算法训练营第四十九天 | 121. 买卖股票的最佳时机,122.买卖股票的最佳时机II

目录 121. 买卖股票的最佳时机 122.买卖股票的最佳时机II 121. 买卖股票的最佳时机 题目链接&#xff1a;121. 买卖股票的最佳时机 比较好想到的是贪心算法&#xff0c;贪心选择价格更低的日子买入&#xff0c;价格更高的日子卖出。 class Solution { public:int maxProfit(vec…

ROS命令行工具

1、roscore 在使用ROS之前&#xff0c;首先要启动roscore进程。当我们在终端中运行这个命令时&#xff0c;系统就会启动ROS Master、参数服务器和日志节点。在这之后&#xff0c;就可以运行任何其他的ROS程序&#xff0f;节点了。所以可以在一个终端窗口运行roscore指令&#…

【unity实战】基于权重的随机事件(附项目源码)

文章目录 前言开始一、简单的使用二、完善各种事件1. 完善生成金币事件2. 完善生成敌人事件敌人3. 完善生成药水事件 最终效果参考源码完结 前言 随机功能和UnityEvent前面其实我们都已经做过了&#xff0c;但是随机UnityEvent事件要怎么使用呢&#xff1f;这里就来举一个例子…

【Linux驱动开发】编译Android12源码+

编译Android12源码 1. 简单描述2. 准备资料3. 编译Android12 1. 简单描述 基于讯为电子rk3568教程 2. 准备资料 rk_android12.0_sdk_20220720.tar.gz 3. 编译Android12 解压 tar -vxf rk_android12.0_sdk_20220720.tar.gz设置屏幕配置 rk_android12.0_sdk/kernel-4.19/ar…

关于src别名的配置之tsconfig.json配置

tsconfig.json {"compilerOptions": {"baseUrl": "./", // 解析非相对模块的基地址&#xff0c;默认是当前目录"paths": { //路径映射&#xff0c;相对于baseUrl"/*": ["src/*"] }} } ① "baseUrl": &…