Rabbitmq之高级特性——实现消费端限流&NACK重回队列_dhtqr-程序员宅基地

  如果是高并发下,rabbitmq服务器上收到成千上万条消息,那么当打开消费端时,这些消息必定喷涌而来,导致消费端消费不过来甚至挂掉都有可能。

在非自动确认的模式下,可以采用限流模式,rabbitmq 提供了服务质量保障qos机制来控制一次消费消息数量。

下面直接上代码:

生产端:

 1 package com.zxy.demo.rabbitmq;
 2 
 3 import java.io.IOException;
 4 import java.util.concurrent.TimeoutException;
 5 
 6 import com.rabbitmq.client.AMQP;
 7 import com.rabbitmq.client.Channel;
 8 import com.rabbitmq.client.ConfirmListener;
 9 import com.rabbitmq.client.Connection;
10 import com.rabbitmq.client.ConnectionFactory;
11 import com.rabbitmq.client.ReturnListener;
12 import com.rabbitmq.client.AMQP.BasicProperties;
13 
14 public class Producter {
15 
16     public static void main(String[] args) throws IOException, TimeoutException {
17         // TODO Auto-generated method stub
18         ConnectionFactory factory = new ConnectionFactory();
19         factory.setHost("192.168.10.110");
20         factory.setPort(5672);
21         factory.setUsername("guest");
22         factory.setPassword("guest");
23         factory.setVirtualHost("/");
24         Connection conn = factory.newConnection();
25         Channel channel = conn.createChannel();
26         String exchange001 = "exchange_001";
27         String queue001 = "queue_001";
28         String routingkey = "mq.topic";
29         String body = "hello rabbitmq!===============限流策略";
30 //        开启确认模式
31         channel.confirmSelect();
32 //        循环发送多条消息        
33         for(int i = 0 ;i<10;i++){
34         channel.basicPublish(exchange001, routingkey, null, body.getBytes());
35     }
36         
37 //        添加一个返回监听========消息返回模式重要添加
38         channel.addConfirmListener(new ConfirmListener() {
39             
40             @Override
41             public void handleNack(long deliveryTag, boolean multiple) throws IOException {
42                 System.out.println("===========NACK============");
43                 
44             }
45             
46             @Override
47             public void handleAck(long deliveryTag, boolean multiple) throws IOException {
48                 System.out.println("===========ACK============");
49                 
50             }
51         });
52     }
53 
54 }

 

消费端:

 1 package com.zxy.demo.rabbitmq;
 2 
 3 import java.io.IOException;
 4 import java.util.concurrent.TimeoutException;
 5 
 6 import com.rabbitmq.client.Channel;
 7 import com.rabbitmq.client.Connection;
 8 import com.rabbitmq.client.ConnectionFactory;
 9 
10 public class Receiver {
11 
12     public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
13         // TODO Auto-generated method stub
14         ConnectionFactory factory = new ConnectionFactory();
15         factory.setHost("192.168.10.110");
16         factory.setPort(5672);
17         factory.setUsername("guest");
18         factory.setPassword("guest");
19         factory.setVirtualHost("/");
20         Connection conn = factory.newConnection();
21         Channel channel = conn.createChannel();
22         String exchange001 = "exchange_001";
23         String queue001 = "queue_001";
24         String routingkey = "mq.*";
25         channel.exchangeDeclare(exchange001, "topic", true, false, null);
26         channel.queueDeclare(queue001, true, false, false, null);
27         channel.queueBind(queue001, exchange001, routingkey);
28 //        设置限流策略
29 //        channel.basicQos(获取消息最大数[0-无限制], 依次获取数量, 作用域[true作用于整个channel,false作用于具体消费者]);
30         channel.basicQos(0, 2, false);
31 //        自定义消费者
32         MyConsumer myConsumer = new MyConsumer(channel);
33 //        进行消费,签收模式一定要为手动签收
34         Thread.sleep(3000);
35         channel.basicConsume(queue001, false, myConsumer);
36     }
37 
38 }

 

自定义消费者:

 1 package com.zxy.demo.rabbitmq;
 2 
 3 import java.io.IOException;
 4 
 5 import com.rabbitmq.client.AMQP.BasicProperties;
 6 import com.rabbitmq.client.Channel;
 7 import com.rabbitmq.client.DefaultConsumer;
 8 import com.rabbitmq.client.Envelope;
 9 
10 /**
11  * 可以继承,可以实现,实现的话要覆写的方法比较多,所以这里用了继承
12  *
13  */
14 public class MyConsumer extends DefaultConsumer{
15     private Channel channel;
16     public MyConsumer(Channel channel) {
17         super(channel);
18         // TODO Auto-generated constructor stub
19         this.channel=channel;
20     }
21 
22     @Override
23     public void handleDelivery(String consumerTag, Envelope envelope, BasicProperties properties, byte[] body)
24             throws IOException {
25         System.out.println("消费标签:"+consumerTag);
26         System.out.println("envelope.getDeliveryTag():==="+envelope.getDeliveryTag());
27         System.out.println("envelope.getExchange():==="+envelope.getExchange());
28         System.out.println("envelope.getRoutingKey():==="+envelope.getRoutingKey());
29         System.out.println("body:==="+new String(body));
30 //        手动签收,一定要有消费者签收,如果没有如下代码,则限流模式下,仅能打印出来channel.basicQos(0, 2, false);第二参数的2条信息
31         channel.basicAck(envelope.getDeliveryTag(), false);
32     }
33     
34 
35 }

 重回队列模式,是当投递消息失败时,让该消息重新回到队列的模式,该模式需要手动签收,并需要在消费者中进行判断,调用重回队列的确认模式

代码如下

生产端:

 1 package com.zxy.demo.rabbitmq;
 2 
 3 import java.io.IOException;
 4 import java.util.HashMap;
 5 import java.util.Map;
 6 import java.util.concurrent.TimeoutException;
 7 
 8 import org.springframework.amqp.core.Message;
 9 
10 import com.rabbitmq.client.AMQP.BasicProperties;
11 import com.rabbitmq.client.Channel;
12 import com.rabbitmq.client.ConfirmListener;
13 import com.rabbitmq.client.Connection;
14 import com.rabbitmq.client.ConnectionFactory;
15 import com.rabbitmq.client.ReturnListener;
16 
17 public class Producter {
18 
19     public static void main(String[] args) throws IOException, TimeoutException {
20         // TODO Auto-generated method stub
21         ConnectionFactory factory = new ConnectionFactory();
22         factory.setHost("192.168.10.110");
23         factory.setPort(5672);
24         factory.setUsername("guest");
25         factory.setPassword("guest");
26         factory.setVirtualHost("/");
27         Connection conn = factory.newConnection();
28         Channel channel = conn.createChannel();
29         String exchange001 = "exchange_001";
30         String queue001 = "queue_001";
31         String routingkey = "mq.topic";
32         
33 //        循环发送多条消息        
34         for(int i = 0 ;i<5;i++){
35             String body = "hello rabbitmq!===============ACK&重回队列,第"+i+"条";
36             Map<String,Object> head = new HashMap<>();
37             head.put("n", i);
38             BasicProperties properties = new BasicProperties(null, "utf-8", head, 2, 1, null, null, null, null, null, null, null, null, null);
39             
40         channel.basicPublish(exchange001, routingkey, properties, body.getBytes());
41     }
42         
43     }
44 
45 }

 

消费端:

 1 package com.zxy.demo.rabbitmq;
 2 
 3 import java.io.IOException;
 4 import java.util.concurrent.TimeoutException;
 5 
 6 import com.rabbitmq.client.Channel;
 7 import com.rabbitmq.client.Connection;
 8 import com.rabbitmq.client.ConnectionFactory;
 9 
10 public class Receiver {
11 
12     public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
13         // TODO Auto-generated method stub
14         ConnectionFactory factory = new ConnectionFactory();
15         factory.setHost("192.168.10.110");
16         factory.setPort(5672);
17         factory.setUsername("guest");
18         factory.setPassword("guest");
19         factory.setVirtualHost("/");
20         Connection conn = factory.newConnection();
21         Channel channel = conn.createChannel();
22         String exchange001 = "exchange_001";
23         String queue001 = "queue_001";
24         String routingkey = "mq.*";
25         channel.exchangeDeclare(exchange001, "topic", true, false, null);
26         channel.queueDeclare(queue001, true, false, false, null);
27         channel.queueBind(queue001, exchange001, routingkey);
28 //        自定义消费者
29         MyConsumer myConsumer = new MyConsumer(channel);
30 //        进行消费,签收模式一定要为手动签收
31         channel.basicConsume(queue001, false, myConsumer);
32     }
33 
34 }

 

自定义消费者:

 1 package com.zxy.demo.rabbitmq;
 2 
 3 import java.io.IOException;
 4 
 5 import com.rabbitmq.client.AMQP.BasicProperties;
 6 import com.rabbitmq.client.Channel;
 7 import com.rabbitmq.client.DefaultConsumer;
 8 import com.rabbitmq.client.Envelope;
 9 
10 /**
11  * 可以继承,可以实现,实现的话要覆写的方法比较多,所以这里用了继承
12  *
13  */
14 public class MyConsumer extends DefaultConsumer{
15     private Channel channel;
16     public MyConsumer(Channel channel) {
17         super(channel);
18         // TODO Auto-generated constructor stub
19         this.channel=channel;
20     }
21 
22     @Override
23     public void handleDelivery(String consumerTag, Envelope envelope, BasicProperties properties, byte[] body)
24             throws IOException {
25         System.out.println("消费标签:"+consumerTag);
26         System.out.println("envelope.getDeliveryTag():==="+envelope.getDeliveryTag());
27         System.out.println("envelope.getExchange():==="+envelope.getExchange());
28         System.out.println("envelope.getRoutingKey():==="+envelope.getRoutingKey());
29         System.out.println("body:==="+new String(body));
30         System.out.println("===================休眠以便查看===============");
31         try {
32             Thread.sleep(2000);
33         } catch (InterruptedException e) {
34             // TODO Auto-generated catch block
35             e.printStackTrace();
36         }
37 //        手动签收
38         Integer i = (Integer) properties.getHeaders().get("n");
39         System.out.println("iiiiiiiiiiiiiiiii======================================================"+i);
40         if(i==1) {
41             channel.basicNack(envelope.getDeliveryTag(),false, true);//第三个参数为是否重返队列
42         }else {
43             channel.basicAck(envelope.getDeliveryTag(), false);    
44         }
45     }
46     
47 
48 }

下面是重回队列执行结果,可以看到当消费完后第一条不断的被扔回队列然后消费再扔回。

 1 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
 2 envelope.getDeliveryTag():===1
 3 envelope.getExchange():===exchange_001
 4 envelope.getRoutingKey():===mq.topic
 5 body:===hello rabbitmq!===============ACK&重回队列,第0条
 6 ===================休眠以便查看===============
 7 iiiiiiiiiiiiiiiii======================================================0
 8 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
 9 envelope.getDeliveryTag():===2
10 envelope.getExchange():===exchange_001
11 envelope.getRoutingKey():===mq.topic
12 body:===hello rabbitmq!===============ACK&重回队列,第1条
13 ===================休眠以便查看===============
14 iiiiiiiiiiiiiiiii======================================================1
15 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
16 envelope.getDeliveryTag():===3
17 envelope.getExchange():===exchange_001
18 envelope.getRoutingKey():===mq.topic
19 body:===hello rabbitmq!===============ACK&重回队列,第2条
20 ===================休眠以便查看===============
21 iiiiiiiiiiiiiiiii======================================================2
22 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
23 envelope.getDeliveryTag():===4
24 envelope.getExchange():===exchange_001
25 envelope.getRoutingKey():===mq.topic
26 body:===hello rabbitmq!===============ACK&重回队列,第3条
27 ===================休眠以便查看===============
28 iiiiiiiiiiiiiiiii======================================================3
29 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
30 envelope.getDeliveryTag():===5
31 envelope.getExchange():===exchange_001
32 envelope.getRoutingKey():===mq.topic
33 body:===hello rabbitmq!===============ACK&重回队列,第4条
34 ===================休眠以便查看===============
35 iiiiiiiiiiiiiiiii======================================================4
36 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
37 envelope.getDeliveryTag():===6
38 envelope.getExchange():===exchange_001
39 envelope.getRoutingKey():===mq.topic
40 body:===hello rabbitmq!===============ACK&重回队列,第1条
41 ===================休眠以便查看===============
42 iiiiiiiiiiiiiiiii======================================================1
43 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
44 envelope.getDeliveryTag():===7
45 envelope.getExchange():===exchange_001
46 envelope.getRoutingKey():===mq.topic
47 body:===hello rabbitmq!===============ACK&重回队列,第1条
48 ===================休眠以便查看===============
49 iiiiiiiiiiiiiiiii======================================================1
50 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
51 envelope.getDeliveryTag():===8
52 envelope.getExchange():===exchange_001
53 envelope.getRoutingKey():===mq.topic
54 body:===hello rabbitmq!===============ACK&重回队列,第1条
55 ===================休眠以便查看===============
56 iiiiiiiiiiiiiiiii======================================================1
57 消费标签:amq.ctag-eeeaa9jgrJ3tDhtQR2XReg
58 envelope.getDeliveryTag():===9
59 envelope.getExchange():===exchange_001
60 envelope.getRoutingKey():===mq.topic
61 body:===hello rabbitmq!===============ACK&重回队列,第1条
62 ===================休眠以便查看===============
View Code

 

版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。
本文链接:https://blog.csdn.net/zuixiaoyao_001/article/details/82599123

智能推荐

5个超厉害的资源搜索网站,每一款都可以让你的资源满满!_最全资源搜索引擎-程序员宅基地

文章浏览阅读1.6w次,点赞8次,收藏41次。生活中我们无时不刻不都要在网站搜索资源,但就是缺少一个趁手的资源搜索网站,如果有一个比较好的资源搜索网站可以帮助我们节省一大半时间!今天小编在这里为大家分享5款超厉害的资源搜索网站,每一款都可以让你的资源丰富精彩!网盘传奇一款最有效的网盘资源搜索网站你还在为找网站里面的资源而烦恼找不到什么合适的工具而烦恼吗?这款网站传奇网站汇聚了4853w个资源,并且它每一天都会持续更新资源;..._最全资源搜索引擎

Book类的设计(Java)_6-1 book类的设计java-程序员宅基地

文章浏览阅读4.5k次,点赞5次,收藏18次。阅读测试程序,设计一个Book类。函数接口定义:class Book{}该类有 四个私有属性 分别是 书籍名称、 价格、 作者、 出版年份,以及相应的set 与get方法;该类有一个含有四个参数的构造方法,这四个参数依次是 书籍名称、 价格、 作者、 出版年份 。裁判测试程序样例:import java.util.*;public class Main { public static void main(String[] args) { List <Book>_6-1 book类的设计java

基于微信小程序的校园导航小程序设计与实现_校园导航微信小程序系统的设计与实现-程序员宅基地

文章浏览阅读613次,点赞28次,收藏27次。相比于以前的传统手工管理方式,智能化的管理方式可以大幅降低学校的运营人员成本,实现了校园导航的标准化、制度化、程序化的管理,有效地防止了校园导航的随意管理,提高了信息的处理速度和精确度,能够及时、准确地查询和修正建筑速看等信息。课题主要采用微信小程序、SpringBoot架构技术,前端以小程序页面呈现给学生,结合后台java语言使页面更加完善,后台使用MySQL数据库进行数据存储。微信小程序主要包括学生信息、校园简介、建筑速看、系统信息等功能,从而实现智能化的管理方式,提高工作效率。

有状态和无状态登录

传统上用户登陆状态会以 Session 的形式保存在服务器上,而 Session ID 则保存在前端的 Cookie 中;而使用 JWT 以后,用户的认证信息将会以 Token 的形式保存在前端,服务器不需要保存任何的用户状态,这也就是为什么 JWT 被称为无状态登陆的原因,无状态登陆最大的优势就是完美支持分布式部署,可以使用一个 Token 发送给不同的服务器,而所有的服务器都会返回同样的结果。有状态和无状态最大的区别就是服务端会不会保存客户端的信息。

九大角度全方位对比Android、iOS开发_ios 开发角度-程序员宅基地

文章浏览阅读784次。发表于10小时前| 2674次阅读| 来源TechCrunch| 19 条评论| 作者Jon EvansiOSAndroid应用开发产品编程语言JavaObjective-C摘要:即便Android市场份额已经超过80%,对于开发者来说,使用哪一个平台做开发仍然很难选择。本文从开发环境、配置、UX设计、语言、API、网络、分享、碎片化、发布等九个方面把Android和iOS_ios 开发角度

搜索引擎的发展历史

搜索引擎的发展历史可以追溯到20世纪90年代初,随着互联网的快速发展和信息量的急剧增加,人们开始感受到了获取和管理信息的挑战。这些阶段展示了搜索引擎在技术和商业模式上的不断演进,以满足用户对信息获取的不断增长的需求。

随便推点

控制对象的特性_控制对象特性-程序员宅基地

文章浏览阅读990次。对象特性是指控制对象的输出参数和输入参数之间的相互作用规律。放大系数K描述控制对象特性的静态特性参数。它的意义是:输出量的变化量和输入量的变化量之比。时间常数T当输入量发生变化后,所引起输出量变化的快慢。(动态参数) ..._控制对象特性

FRP搭建内网穿透(亲测有效)_locyanfrp-程序员宅基地

文章浏览阅读5.7w次,点赞50次,收藏276次。FRP搭建内网穿透1.概述:frp可以通过有公网IP的的服务器将内网的主机暴露给互联网,从而实现通过外网能直接访问到内网主机;frp有服务端和客户端,服务端需要装在有公网ip的服务器上,客户端装在内网主机上。2.简单的图解:3.准备工作:1.一个域名(www.test.xyz)2.一台有公网IP的服务器(阿里云、腾讯云等都行)3.一台内网主机4.下载frp,选择适合的版本下载解压如下:我这里服务器端和客户端都放在了/usr/local/frp/目录下4.执行命令# 服务器端给执_locyanfrp

UVA 12534 - Binary Matrix 2 (网络流‘最小费用最大流’ZKW)_uva12534-程序员宅基地

文章浏览阅读687次。题目:http://acm.hust.edu.cn/vjudge/contest/view.action?cid=93745#problem/A题意:给出r*c的01矩阵,可以翻转格子使得0表成1,1变成0,求出最小的步数使得每一行中1的个数相等,每一列中1的个数相等。思路:网络流。容量可以保证每一行和每一列的1的个数相等,费用可以算出最小步数。行向列建边,如果该格子是_uva12534

免费SSL证书_csdn alphassl免费申请-程序员宅基地

文章浏览阅读504次。1、Let's Encrypt 90天,支持泛域名2、Buypass:https://www.buypass.com/ssl/resources/go-ssl-technical-specification6个月,单域名3、AlwaysOnSLL:https://alwaysonssl.com/ 1年,单域名 可参考蜗牛(wn789)4、TrustAsia5、Alpha..._csdn alphassl免费申请

测试算法的性能(以选择排序为例)_算法性能测试-程序员宅基地

文章浏览阅读1.6k次。测试算法的性能 很多时候我们需要对算法的性能进行测试,最简单的方式是看算法在特定的数据集上的执行时间,简单的测试算法性能的函数实现见testSort()。【思想】:用clock_t计算某排序算法所需的时间,(endTime - startTime)/ CLOCKS_PER_SEC来表示执行了多少秒。【关于宏CLOCKS_PER_SEC】:以下摘自百度百科,“CLOCKS_PE_算法性能测试

Lane Detection_lanedetectionlite-程序员宅基地

文章浏览阅读1.2k次。fromhttps://towardsdatascience.com/finding-lane-lines-simple-pipeline-for-lane-detection-d02b62e7572bIdentifying lanes of the road is very common task that human driver performs. This is important ..._lanedetectionlite

推荐文章

热门文章

相关标签