发送单个消息的生产者和接收消息并打印出来的消费者。
在下图中,“ P”是我们的生产者,“ C”是我们的消费者。中间的框是一个队列-RabbitMQ 代表使用者保留的消息缓冲区
<dependencies> <!--rabbitmq依赖客户端--> <dependency> <groupId>com.rabbitmq</groupId> <artifactId>amqp-client</artifactId> <version>5.9.0</version> </dependency> <!--操作文件流依赖--> <dependency> <groupId>commons-io</groupId> <artifactId>commons-io</artifactId> <version>2.6</version> </dependency> </dependencies>
package com.study.rabbitmq.one; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import java.io.IOException; import java.util.concurrent.TimeoutException; public class Producer { //队列名称 public static final String QUEUE_NAME = "hello"; //发消息 public static void main(String[] args) throws IOException, TimeoutException { //创建一个连接工厂 ConnectionFactory factory = new ConnectionFactory(); //工厂IP,连接RabbitMQ队列 factory.setHost("192.168.137.4"); //连接端口号 factory.setPort(5672); //用户名 factory.setUsername("admin"); //密码 factory.setPassword("123"); //创建连接 Connection connection = factory.newConnection(); //获取信道 Channel channel = connection.createChannel(); /* * 生成一个队列 * 1.队列名称 * 2.队列里面的消息是否持久化(存储在磁盘),默认情况消息存储在内存中 * 3.该队列是否只供一个消费者进行消费,是否进行消息共享。true可以多个消费者消费,false只能一个消费者消费 * 4.最后一个消费者端开链接以后该队列是否自动删除 true自动删除 false不自动删除 * */ channel.queueDeclare(QUEUE_NAME,false,false,false,null); //发消息 String message = "hello world"; /* *发送一次消费 * 1.发送到哪个交换机 * 2.路由的key值是哪个 本次是队列的名称 * 3.其他参数信息 * 4.发送消息的消息体 * */ channel.basicPublish("",QUEUE_NAME,null,message.getBytes()); System.out.println("消息发送完毕"); } }
package com.study.rabbitmq.one; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; public class Consumer { //队列名称 public static final String QUEUE_NAME = "hello"; //接收消息 public static void main(String[] args) throws IOException, TimeoutException { //创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.137.4"); factory.setPort(5672); factory.setUsername("admin"); factory.setPassword("123"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); //声明接收消息 DeliverCallback deliverCallback = (consumerTag,message) -> { System.out.println(new String(message.getBody())); }; //取消消息时的回调 CancelCallback cancelCallback = consumerTag -> { System.out.println("消息消费被中断"); }; /* * 消费者消费消息 * 1.消费哪个队列 * 2.消费成功之后是否要自动应答 true自动应答 false手动应答 * 3.消费者未成功消费的回调 * 4.消费者取消消费的回调 * */ channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback); } }
工作队列(又称任务队列)的主要思想是避免立即执行资源密集型任务,而不得不等待它完成。相反我们安排任务在之后执行。我们把任务封装为消息并将其发送到队列。在后台运行的工作进程将弹出任务并最终执行作业。当有多个工作线程时,这些工作线程将一起处理这些任务。
轮询分发消息
一个生产者发送消息,由多个工作线程(消费者)轮询接收
在这个案例中我们会启动两个工作线程,一个消息发送线程,我们来看看他们两个工作线程是如何工作的。
一、编写工具类,提取重复代码
package com.study.rabbitmq.utils; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; public class RabbitMQUtils { //得到一个连接的 channel public static Channel getChannel() throws Exception{ //创建一个连接工厂 ConnectionFactory factory = new ConnectionFactory(); //工厂IP,连接RabbitMQ队列 factory.setHost("192.168.137.4"); //连接端口号 factory.setPort(5672); //用户名 factory.setUsername("admin"); //密码 factory.setPassword("123"); //创建连接 Connection connection = factory.newConnection(); //获取信道 Channel channel = connection.createChannel(); return channel; } }
二、编写消息发送线程,在控制台输入发送的消息
package com.study.rabbitmq.two; import com.rabbitmq.client.Channel; import com.study.rabbitmq.utils.RabbitMQUtils; import java.util.Scanner; //生产者 public class Task01 { //队列名称 public static final String QUEUE_NAME = "hello02"; //发送大量消息 public static void main(String[] args) throws Exception { Channel channel = RabbitMQUtils.getChannel(); //声明队列 channel.queueDeclare(QUEUE_NAME,false,false,false,null); //从控制台当中接收信息 Scanner scanner = new Scanner(System.in); while (scanner.hasNext()){ String message = scanner.next(); /* *发送一次消费 * 1.发送到哪个交换机 * 2.路由的key值是哪个 本次是队列的名称 * 3.其他参数信息 * 4.发送消息的消息体 * */ channel.basicPublish("",QUEUE_NAME,null,message.getBytes()); System.out.println("发送消息完成"+message); } } }
三、编写两个接收消息的工作线程
package com.study.rabbitmq.two; import com.rabbitmq.client.CancelCallback; import com.rabbitmq.client.Channel; import com.rabbitmq.client.DeliverCallback; import com.study.rabbitmq.utils.RabbitMQUtils; //工作线程 public class Worker01 { //队列名称 public static final String QUEUE_NAME = "hello02"; public static void main(String[] args) throws Exception { Channel channel = RabbitMQUtils.getChannel(); //接收消息 DeliverCallback deliverCallback = (consumerTag,message) ->{ System.out.println("接收到的消息"+new String(message.getBody())); }; //消息接收被取消时,执行下面的内容 CancelCallback cancelCallback =(consumerTag) -> { System.out.println(consumerTag+":消息取消消费接口回调逻辑"); }; System.out.println("C1等待接收消息"); //消息接收 /* * 消费者消费消息 * 1.消费哪个队列 * 2.消费成功之后是否要自动应答 true自动应答 false手动应答 * 3.消费者未成功消费的回调 * 4.消费者取消消费的回调 * */ channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback); } }
package com.study.rabbitmq.two; import com.rabbitmq.client.CancelCallback; import com.rabbitmq.client.Channel; import com.rabbitmq.client.DeliverCallback; import com.study.rabbitmq.utils.RabbitMQUtils; //工作线程 public class Worker02 { //队列名称 public static final String QUEUE_NAME = "hello02"; public static void main(String[] args) throws Exception { Channel channel = RabbitMQUtils.getChannel(); //接收消息 DeliverCallback deliverCallback = (consumerTag, message) ->{ System.out.println("接收到的消息"+new String(message.getBody())); }; //消息接收被取消时,执行下面的内容 CancelCallback cancelCallback =(consumerTag) -> { System.out.println(consumerTag+":消息取消消费接口回调逻辑"); }; System.out.println("C2等待接收消息"); //消息接收 /* * 消费者消费消息 * 1.消费哪个队列 * 2.消费成功之后是否要自动应答 true自动应答 false手动应答 * 3.消费者未成功消费的回调 * 4.消费者取消消费的回调 * */ channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback); } }
四、先启动消息发送线程创建hello02信道,再启动两个接收消息的工作线程
五、在消息发送线程控制台输入以下内容
六、查看两个工作线程分别接收到的消息
哪个线程先启动,哪个最先接收消息!