RabbitMQ学习:RabbitMQ的六种工作模式之简单和工作模式(三)

上一篇:RabbitMQ学习:RabbitMQ的基本概念及RabbitMQ使用场景(二) --- https://my.oschina.net/u/4115134/blog/3223371


RabbitMQ的六种工作模式

首先开启虚拟机上的rabbitmq服务器

1# 启动服务 2systemctl start rabbitmq-server

一、简单模式

RabbitMQ是一个消息中间件,你可以想象它是一个邮局。当你把信件放到邮箱里时,能够确信邮递员会正确地递送你的信件。RabbitMq就是一个邮箱、一个邮局和一个邮递员。

  • 发送消息的程序是生产者

  • 队列就代表一个邮箱。虽然消息会流经RbbitMQ和你的应用程序,但消息只能被存储在队列里。队列存储空间只受服务器内存和磁盘限制,它本质上是一个大的消息缓冲区。多个生产者可以向同一个队列发送消息,多个消费者也可以从同一个队列接收消息.

  • 消费者等待从队列接收消息

创建Rabbitmq-demo 的测试项

1、pom.xml

添加 slf4j 依赖, 和 rabbitmq amqp 依赖

1<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> 2 <modelVersion>4.0.0</modelVersion> 3 <groupId>com.qile</groupId> 4 <artifactId>rabbitmq</artifactId> 5 <version>0.0.1-SNAPSHOT</version> 6 7 <dependencies> 8 <dependency> 9 <groupId>com.rabbitmq</groupId> 10 <artifactId>amqp-client</artifactId> 11 <version>5.4.3</version> 12 </dependency> 13 <dependency> 14 <groupId>org.slf4j</groupId> 15 <artifactId>slf4j-api</artifactId> 16 <version>1.8.0-alpha2</version> 17 </dependency> 18 <dependency> 19 <groupId>org.slf4j</groupId> 20 <artifactId>slf4j-log4j12</artifactId> 21 <version>1.8.0-alpha2</version> 22 </dependency> 23 </dependencies> 24 25 <build> 26 <plugins> 27 <plugin> 28 <groupId>org.apache.maven.plugins</groupId> 29 <artifactId>maven-compiler-plugin</artifactId> 30 <version>3.8.0</version> 31 <configuration> 32 <source>1.8</source> 33 <target>1.8</target> 34 </configuration> 35 </plugin> 36 </plugins> 37 </build> 38 39</project>

2. 生产者发送消息--HelloWorld

1package rabbitmq.simple; 2 3import java.io.IOException; 4import java.util.concurrent.TimeoutException; 5 6import com.rabbitmq.client.Channel; 7import com.rabbitmq.client.Connection; 8import com.rabbitmq.client.ConnectionFactory; 9 10public class Test1 { 11 12 public static void main(String[] args) throws IOException, TimeoutException { 13 14 /** 15 * 1. 建立连接 16 * 2. 创建队列:helloworld 17 * 3. 向队列发送数据 18 */ 19 ConnectionFactory f = new ConnectionFactory(); 20 f.setHost("192.168.64.140"); 21 f.setPort(5672); 22 f.setUsername("admin"); 23 f.setPassword("admin"); 24 25 /* 26 * 与rabbitmq服务器建立连接, 27 * rabbitmq服务器端使用的是nio,会复用tcp连接, 28 * 并开辟多个信道与客户端通信 29 * 以减轻服务器端建立连接的开销 30 */ 31 Connection con = f.newConnection(); 32 //创建通道 33 Channel c = con.createChannel(); 34 35 /* 36 * 声明队列,会在rabbitmq中创建一个队列 37 * 如果已经创建过该队列,就不能再使用其他参数来创建 38 * 39 * 参数含义: 40 * -queue: 队列名称 41 * -durable: 队列持久化,true表示RabbitMQ重启后队列仍存在 42 * -exclusive: 排他,true表示限制仅当前连接可用 43 * -autoDelete: 当最后一个消费者断开后,是否删除队列 44 * -arguments: 其他参数 45 */ 46 c.queueDeclare("helloworld",false,false,false,null); 47 48 /* 49 * 发布消息 50 * 这里把消息向默认交换机发送. 51 * 默认交换机隐含与所有队列绑定,routing key即为队列名称 52 * 53 * 参数含义: 54 * -exchange: 交换机名称,空串表示默认交换机"(AMQP default)",不能用 null 55 * -routingKey: 对于默认交换机,路由键就是目标队列名称 56 * -props: 其他参数,例如头信息 57 * -body: 消息内容byte[]数组 58 */ 59 c.basicPublish("","helloworld",null, 60 ("Hello World!" + System.currentTimeMillis()).getBytes()); 61 System.out.println("消息已发出"); 62 63 c.close(); 64 con.close(); 65 } 66}

这时Run as 得到

在rabbitmq客户端有:

之后编写消费者接受消息

3、消费者接收消息

1package rabbitmq.simple; 2 3import java.io.IOException; 4import java.util.concurrent.TimeoutException; 5 6import com.rabbitmq.client.CancelCallback; 7import com.rabbitmq.client.Channel; 8import com.rabbitmq.client.Connection; 9import com.rabbitmq.client.ConnectionFactory; 10import com.rabbitmq.client.DeliverCallback; 11import com.rabbitmq.client.Delivery; 12 13public class Test2 { 14 15 public static void main(String[] args) throws IOException, TimeoutException { 16 17 /** 18 * 1. 建立连接 19 * 2. 创建队列:helloworld 20 * 3. 向队列发送数据 21 */ 22 ConnectionFactory f = new ConnectionFactory(); 23 f.setHost("192.168.64.140"); 24 f.setPort(5672); 25 f.setUsername("admin"); 26 f.setPassword("admin"); 27 Connection con = f.newConnection(); //创建连接 28 Channel c = con.createChannel(); //创建通道 29 30 //定义队列,服务器没有这个队列会创建,若有什么都不做 31 c.queueDeclare("helloworld",false,false,false,null); 32 33 //收到消息后用来处理消息的回调对象 34 DeliverCallback deliverCallback = new DeliverCallback() { 35 @Override 36 public void handle(String consumerTag, Delivery message) throws IOException { 37 byte[] a = message.getBody(); 38 String msg = new String(a); 39 System.out.println("收到" + msg); 40 } 41 }; 42 43 //消费者取消时的回调对象 44 CancelCallback cancelCallback = new CancelCallback() { 45 @Override 46 public void handle(String consumerTag) throws IOException { 47 48 } 49 }; 50 51 //开始消费数据 52 c.basicConsume("helloworld",true,deliverCallback,cancelCallback); 53 } 54 55}

此时,在之前所积累的两条消息将会在你程序运转之时,显示出来,这是再去运转生产者,将会直接显示出发送的数据

二、工作模式

工作队列(即任务队列)背后的主要思想是避免立即执行资源密集型任务,并且必须等待它完成。相反,我们将任务安排在稍后完成。

我们将任务封装为消息并将其发送到队列。后台运行的工作进程将获取任务并最终执行任务。当运行多个消费者时,任务将在它们之间分发

使用任务队列的一个优点是能够轻松地并行工作。如果我们正在积压工作任务,我们可以添加更多工作进程,这样就可以轻松扩展。

1、生产者发送消息

这里模拟耗时任务,发送的消息中,每个点使工作进程暂停一秒钟,例如"Hello…"将花费3秒钟来处理

1package rabbitmq.work; 2 3import java.util.Scanner; 4 5import com.rabbitmq.client.Channel; 6import com.rabbitmq.client.Connection; 7import com.rabbitmq.client.ConnectionFactory; 8 9public class Test1 { 10 public static void main(String[] args) throws Exception { 11 /** 12 * 1. 建立连接 13 * 2. 创建队列:helloworld 14 * 3. 向队列发送数据 15 */ 16 ConnectionFactory f = new ConnectionFactory(); 17 f.setHost("192.168.64.140"); 18 f.setPort(5672); 19 f.setUsername("admin"); 20 f.setPassword("admin"); 21 22 Connection c = f.newConnection(); //创建连接 23 Channel ch = c.createChannel(); //创建通道 24 //参数:queue,durable,exclusive,autoDelete,arguments 25 ch.queueDeclare("helloworld", false,false,false,null); 26 27 /** 28 * 模拟耗时消息 29 * 发送的字符串中,有一个点字符,消费者处理的时候就暂停1秒 30 */ 31 //循环输入消息发送到rabbitmq 32 while (true) { 33 System.out.print("输入消息: "); 34 String msg = new Scanner(System.in).nextLine(); 35 //如果输入的是"exit"则结束生产者进程 36 if ("exit".equals(msg)) { 37 break; 38 } 39 //参数:exchage,routingKey,props,body 40 ch.basicPublish("", "helloworld", null, msg.getBytes()); 41 System.out.println("消息已发送: "+msg); 42 } 43 44 c.close(); 45 } 46}

2、消费者接收消息

1package rabbitmq.work; 2 3import java.io.IOException; 4import java.util.concurrent.TimeoutException; 5 6import com.rabbitmq.client.CancelCallback; 7import com.rabbitmq.client.Channel; 8import com.rabbitmq.client.Connection; 9import com.rabbitmq.client.ConnectionFactory; 10import com.rabbitmq.client.DeliverCallback; 11import com.rabbitmq.client.Delivery; 12 13public class Test2 { 14 public static void main(String[] args) throws Exception { 15 16 /** 17 * 1. 建立连接 18 * 2. 创建队列:helloworld 19 * 3. 向队列发送数据 20 */ 21 ConnectionFactory f = new ConnectionFactory(); 22 f.setHost("192.168.64.140"); 23 f.setUsername("admin"); 24 f.setPassword("admin"); 25 Connection c = f.newConnection(); //创建连接 26 Channel ch = c.createChannel(); //创建通道 27 28 ch.queueDeclare("helloworld",false,false,false,null); 29 System.out.println("等待接收数据"); 30 31 //收到消息后用来处理消息的回调对象 32 DeliverCallback callback = new DeliverCallback() { 33 @Override 34 public void handle(String consumerTag, Delivery message) throws IOException { 35 String msg = new String(message.getBody(), "UTF-8"); 36 System.out.println("收到: "+msg); 37 38 //遍历字符串中的字符,每个点使进程暂停一秒 39 for (int i = 0; i < msg.length(); i++) { 40 if (msg.charAt(i)=='.') { 41 try { 42 Thread.sleep(1000); 43 } catch (InterruptedException e) { 44 } 45 } 46 } 47 System.out.println("处理结束"); 48 } 49 }; 50 51 //消费者取消时的回调对象 52 CancelCallback cancel = new CancelCallback() { 53 @Override 54 public void handle(String consumerTag) throws IOException { 55 } 56 }; 57 58 ch.basicConsume("helloworld", true, callback, cancel); 59 } 60}

3.运行测试

运行:

  • 一个生产者
  • 两个消费者

生产者发送多条消息
如:1,2,3,4,5,...两个消费者分别收到:

  • 消费者一:1,3,5,...
  • 消费者二:2,4,...

rabbtimq在所有消费者中轮询分布消息,把消息均匀发送给所有消费者。

4.消息确认

一个消费者接收消息后,在消息没有完全处理完时就挂掉了,那么这时会发生什么呢?

就现在的代码来说,rabbitmq把消息发送给消费者后,会立即删除消息,那么消费者挂掉后,它没来得及处理的消息就会丢失

1 如果生产者发送以下消息: 2 3 14 5 2 6 7 3 8 9 4 10 11 5 12 13 两个消费者分别收到: 14 15 消费者一: 1, 3, 5 16 消费者二: 2, 4 17 18 当消费者一收到所有消息后,要话费7秒时间来处理第一条消息,这期间如果关闭该消费者,那么1未处理完成,3,5则没有被处理

我们并不想丢失任何消息, 如果一个消费者挂掉,我们想把它的任务消息派发给其他消费者

为了确保消息不会丢失,rabbitmq支持消息确认(回执)。当一个消息被消费者接收到并且执行完成后,消费者会发送一个ack (acknowledgment) 给rabbitmq服务器, 告诉他我已经执行完成了,你可以把这条消息删除了。

如果一个消费者没有返回消息确认就挂掉了(信道关闭,连接关闭或者TCP链接丢失),rabbitmq就会明白,这个消息没有被处理完成rabbitmq就会把这条消息重新放入队列,如果在这时有其他的消费者在线,那么rabbitmq就会迅速的把这条消息传递给其他的消费者,这样就确保了没有消息会丢失。

这里不存在消息超时, rabbitmq只在消费者挂掉时重新分派消息, 即使消费者花非常久的时间来处理消息也可以

手动消息确认默认是开启的,前面的例子我们通过autoAck=ture把它关闭了。我们现在要把它设置为false,然后工作进程处理完意向任务时,发送一个消息确认(回执)。

1package rabbitmq.work; 2 3import java.io.IOException; 4import java.util.concurrent.TimeoutException; 5 6import com.rabbitmq.client.CancelCallback; 7import com.rabbitmq.client.Channel; 8import com.rabbitmq.client.Connection; 9import com.rabbitmq.client.ConnectionFactory; 10import com.rabbitmq.client.DeliverCallback; 11import com.rabbitmq.client.Delivery; 12 13public class Test2 { 14 public static void main(String[] args) throws Exception { 15 16 /** 17 * 1. 建立连接 18 * 2. 创建队列:helloworld 19 * 3. 向队列发送数据 20 */ 21 ConnectionFactory f = new ConnectionFactory(); 22 f.setHost("192.168.64.140"); 23 f.setUsername("admin"); 24 f.setPassword("admin"); 25 Connection c = f.newConnection(); //创建连接 26 Channel ch = c.createChannel(); //创建通道 27 28 //声明队列 29 ch.queueDeclare("helloworld",false,false,false,null); 30 System.out.println("等待接收数据"); 31 32 //收到消息后用来处理消息的回调对象 33 DeliverCallback callback = new DeliverCallback() { 34 @Override 35 public void handle(String consumerTag, Delivery message) throws IOException { 36 String msg = new String(message.getBody(), "UTF-8"); 37 System.out.println("收到: "+msg); 38 for (int i = 0; i < msg.length(); i++) { 39 if (msg.charAt(i)=='.') { 40 try { 41 Thread.sleep(1000); 42 } catch (InterruptedException e) { 43 } 44 } 45 } 46 System.out.println("处理结束"); 47 //发送回执 48 ch.basicAck(message.getEnvelope().getDeliveryTag(), false); 49 } 50 }; 51 52 //消费者取消时的回调对象 53 CancelCallback cancel = new CancelCallback() { 54 @Override 55 public void handle(String consumerTag) throws IOException { 56 } 57 }; 58 59 //autoAck设置为false,则需要手动确认发送回执 60 ch.basicConsume("helloworld", false, callback, cancel); 61 } 62}

使用以上代码,就算杀掉一个正在处理消息的工作进程也不会丢失任何消息,工作进程挂掉之后,没有确认的消息就会被自动重新传递。

忘记确认(ack)是一个常见的错误, 这样后果是很严重的, 由于未确认的消息不会被释放, rabbitmq会吃掉越来越多的内存

可以使用下面命令打印工作队列中未确认消息的数量

rabbitmqctl list_queues name messages_ready messages_unacknowledged

当处理消息时异常中断, 可以选择让消息重回队列重新发送. nack 操作可以是消息重回队列, 可以使用 basicNack() 方法:

1// requeue为true时重回队列, 反之消息被丢弃或被发送到死信队列 2c.basicNack(tag, multiple, requeue)

5.合理地分发

rabbitmq会一次把多个消息分发给消费者, 这样可能造成有的消费者非常繁忙, 而其它消费者空闲. 而rabbitmq对此一无所知, 仍然会均匀的分发消息

我们可以使用 basicQos(1) 方法, 这告诉rabbitmq_一次只向消费者发送一条消息_, 在返回确认回执前, 不要向消费者发送新消息. 而是把消息发给下一个空闲的消费者

6.消息持久化

当rabbitmq关闭时, 我们队列中的消息仍然会丢失, 除非明确要求它不要丢失数据

要求rabbitmq不丢失数据要做如下两点: 把队列和消息都设置为可持久化(durable)

队列设置为可持久化, 可以在定义队列时指定参数durable为true

1//第二个参数是持久化参数durable 2ch.queueDeclare("helloworld", true, false, false, null);

由于之前我们已经定义过队列"hello"是不可持久化的, 对已存在的队列, rabbitmq不允许对其定义不同的参数, 否则会出错, 所以这里我们定义一个不同名字的队列"task_queue"

1//定义一个新的队列,名为 task_queue 2//第二个参数是持久化参数 durable 3ch.queueDeclare("task_queue", true, false, false, null);

生产者和消费者代码都要修改

这样即使rabbitmq重新启动, 队列也不会丢失. 现在我们再设置队列中消息的持久化, 使用MessageProperties.PERSISTENT_TEXT_PLAIN参数

1//第三个参数设置消息持久化 2ch.basicPublish("", "task_queue", 3 MessageProperties.PERSISTENT_TEXT_PLAIN, 4 msg.getBytes());

下面是"工作模式"最终完成的生产者和消费者代码

7.生产者代码

1package rabbitmq.work; 2 3import java.util.Scanner; 4 5import com.rabbitmq.client.Channel; 6import com.rabbitmq.client.Connection; 7import com.rabbitmq.client.ConnectionFactory; 8import com.rabbitmq.client.MessageProperties; 9 10public class Test3 { 11 public static void main(String[] args) throws Exception { 12 ConnectionFactory f = new ConnectionFactory(); 13 f.setHost("192.168.64.140"); 14 f.setPort(5672); 15 f.setUsername("admin"); 16 f.setPassword("admin"); 17 18 Connection c = f.newConnection(); 19 Channel ch = c.createChannel(); 20 21 //第二个参数设置队列持久化 22 ch.queueDeclare("task_queue", true,false,false,null); 23 24 while (true) { 25 System.out.print("输入消息: "); 26 String msg = new Scanner(System.in).nextLine(); 27 if ("exit".equals(msg)) { 28 break; 29 } 30 31 //第三个参数设置消息持久化 32 ch.basicPublish("", "task_queue", MessageProperties.PERSISTENT_TEXT_PLAIN, msg.getBytes("UTF-8")); 33 System.out.println("消息已发送: "+msg); 34 } 35 36 c.close(); 37 } 38}

8.消费者代码

1package rabbitmq.work; 2 3import java.io.IOException; 4import java.util.concurrent.TimeoutException; 5 6import com.rabbitmq.client.CancelCallback; 7import com.rabbitmq.client.Channel; 8import com.rabbitmq.client.Connection; 9import com.rabbitmq.client.ConnectionFactory; 10import com.rabbitmq.client.DeliverCallback; 11import com.rabbitmq.client.Delivery; 12 13public class Test4 { 14 public static void main(String[] args) throws Exception { 15 ConnectionFactory f = new ConnectionFactory(); 16 f.setHost("192.168.64.140"); 17 f.setUsername("admin"); 18 f.setPassword("admin"); 19 Connection c = f.newConnection(); 20 Channel ch = c.createChannel(); 21 22 //定义一个新的队列,名为 task_queue 23 //设定第二个参数是持久化参数 durable为true 24 ch.queueDeclare("task_queue",true,false,false,null); 25 26 System.out.println("等待接收数据"); 27 28 ch.basicQos(1); //一次只接收一条消息 29 30 //收到消息后用来处理消息的回调对象 31 DeliverCallback callback = new DeliverCallback() { 32 @Override 33 public void handle(String consumerTag, Delivery message) throws IOException { 34 String msg = new String(message.getBody(), "UTF-8"); 35 System.out.println("收到: "+msg); 36 for (int i = 0; i < msg.length(); i++) { 37 if (msg.charAt(i)=='.') { 38 try { 39 Thread.sleep(1000); 40 } catch (InterruptedException e) { 41 } 42 } 43 } 44 System.out.println("处理结束"); 45 //发送回执 46 ch.basicAck(message.getEnvelope().getDeliveryTag(), false); 47 } 48 }; 49 50 //消费者取消时的回调对象 51 CancelCallback cancel = new CancelCallback() { 52 @Override 53 public void handle(String consumerTag) throws IOException { 54 } 55 }; 56 57 //autoAck设置为false,则需要手动确认发送回执 58 ch.basicConsume("task_queue", false, callback, cancel); 59 } 60}

9.总结

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

RabbitMQ六种队列模式

前言RabbitMQ六种队列模式简单队列(https://www.oschina.net/action/GoToLink?urlhttps%3A%2F%2Fwww.cnblogs.com%2Fniceyoo%2Fp%2F11448111.html)RabbitMQ六种队列模式工作队列(https://www.oschi

KVM调整cpu和内存

一.修改kvm虚拟机的配置1、virsheditcentos7找到“memory”和“vcpu”标签,将<namecentos7</name<uuid2220a6d1a36a4fbb8523e078b3dfe795</uuid

RabbitMQ学习:RabbitMQ的六种工作模式之简单和工作模式(三) - HelloWorld