6/12-6/13 分布式消息队列(1)
1、消息队列优势
异步,生产者发布后无需等待消费者,可以继续执行自己的的逻辑。
削峰,将大量的请求保存到队列中,服务器依据自身能力逐步处理请求。
2、分布式消息队列优势
数据持久化,可以把消息存储到硬盘,数据不会丢失。
可扩展性,可根据需求动态增加或者减少结点,保持服务的稳定
应用解耦,允许不同框架语言的系统进行数据的传输和读取
(解耦:生产者和消费者之间不直接依赖对方的实现细节(如接口、调用方式、运行状态等),双方只需约定数据格式或协议)
3、rabbitmq基本使用
3.1 helloworld
生产者
注意:
1)远程连接时要打开5672端口网络安全组,也要打开服务器本地防火墙
2)channel可以复用connection的tcp连接,且每个channel也可以有自己的状态。就像一家快递公司(Connection)可以派多个快递员(Channel)同时送货,比每次送货都重新开一家快递公司高效多了。
▼java复制代码import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import java.nio.charset.StandardCharsets; public class Send { private final static String QUEUE_NAME = "hello"; private final static String REMOTE_HOST = "x.x.x.x";// 输入你的服务器 public static void main(String[] argv) throws Exception { // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); factory.setHost(REMOTE_HOST); factory.setUsername("klee"); // 账号 factory.setPassword(""); // 密码 // 创建连接 try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { /** * 队列名 * 是否持久化 * 是否独占发送连接 * 是否自动删除队列 * 其他参数 */ channel.queueDeclare(QUEUE_NAME, false, false, false, null); String message = "Hello World2!"; // 发送消息 channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } }
消费者
▼java复制代码import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.DeliverCallback; import java.nio.charset.StandardCharsets; public class Recv { private final static String QUEUE_NAME = "hello"; private final static String REMOTE_HOST = ""; public static void main(String[] argv) throws Exception { // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); factory.setHost(REMOTE_HOST); factory.setUsername("klee"); factory.setPassword(""); // 从工厂获取一个新的连接 Connection connection = factory.newConnection(); // 从连接中创建一个新的频道 Channel channel = connection.createChannel(); // 创建队列,在该频道上声明我们正在监听的队列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); // 在控制台打印等待接收消息的信息 System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); // 定义了如何处理消息,创建一个新的DeliverCallback来处理接收到的消息 DeliverCallback deliverCallback = (consumerTag, delivery) -> { // 将消息体转换为字符串 String message = new String(delivery.getBody(), StandardCharsets.UTF_8); // 在控制台打印已接收消息的信息 System.out.println(" [x] Received '" + message + "'"); }; // 在频道上开始消费队列中的消息,接收到的消息会传递给deliverCallback来处理,会持续阻塞 channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> { }); } }
3.2 workqueue
生产者
▼java复制代码import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.MessageProperties; import java.util.Scanner; public class NewTask { // 定义要使用的队列名称,这次的队列就改叫multi_queue private static final String TASK_QUEUE_NAME = "task_queue"; public static void main(String[] argv) throws Exception { // 创建一个连接工厂 ConnectionFactory factory = new ConnectionFactory(); // 设置RabbitMQ服务的主机名 factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); // 创建一个新的连接 try (Connection connection = factory.newConnection(); // 创建一个新的频道 Channel channel = connection.createChannel()) { // 声明队列参数,包括队列名称、是否持久化等 channel.queueDeclare(TASK_QUEUE_NAME, false, false, false, null); // 创建一个输入扫描器,用于读取控制台输入 Scanner scanner = new Scanner(System.in); // 使用循环,每当用户在控制台输入一行文本,就将其作为消息发送 while (scanner.hasNext()) { // 读取用户在控制台输入的下一行文本 String message = scanner.nextLine(); // 发布消息到队列,设置消息持久化 channel.basicPublish("", TASK_QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes("UTF-8")); // 输出到控制台,表示消息已发送 System.out.println(" [x] Sent '" + message + "'"); } } } }
多消费者:
注意:
1)channel.basicQos(X); 启动则开启 公平 轮巡,每个消费者最多积压 X 个未确认消息
2)channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); 确认信息,慎用批量确认
3)channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
最后一个参数true为重新入队,false为拒绝或进入死信队列,慎用true小心消息循环
▼java复制代码import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.DeliverCallback; import lombok.extern.slf4j.Slf4j; @Slf4j public class Worker { private static final String TASK_QUEUE_NAME = "task_queue"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); final Connection connection = factory.newConnection(); for(int i=1;i<=2;i++) { final Channel channel = connection.createChannel(); // 声明队列 channel.queueDeclare(TASK_QUEUE_NAME, false, false, false, null); System.out.println("[*] 消费者" + i + "号 Waiting for messages. To exit press CTRL+C"); // 启动则开启公平轮巡,每个消费者最多积压 prefetchCount 个未确认消息 channel.basicQos(1); int finalI = i; // 接收消息 DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [x] Received '" + message + "'"); try { doWork(message, finalI); } catch (InterruptedException e) { log.error("doWork error", e); // 拒绝消息 channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); } finally { System.out.println(" [x]" + message + " Done by 消费者" + finalI); // 确认消息已经被处理 channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } }; // 开启消息接收,建议手动确认 channel.basicConsume(TASK_QUEUE_NAME, false, deliverCallback, consumerTag -> { }); } } private static void doWork(String task,int id) throws InterruptedException { System.out.println("任务:" + task + "已经被消费者id:" + id + "接收"); Thread.sleep(30*1000); // 模拟处理任务需要时间 } }
3.3 fanout交换机
即全体广播,生产者声明一个交换机并指定其类型为fanout,发布消息也指定交换机即可
▼java复制代码package com.zxw.bi_intelligence.rabbitmq.exchanger.fanout; import com.rabbitmq.client.*; import java.util.Scanner; public class EmitLog { private static final String EXCHANGE_NAME = "fanout"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); // 设置RabbitMQ服务的主机名 factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { // 创建一个交换机 channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); // 创建一个输入扫描器,用于读取控制台输入 Scanner scanner = new Scanner(System.in); // 使用循环,每当用户在控制台输入一行文本,就将其作为消息发送 while (scanner.hasNext()) { // 读取用户在控制台输入的下一行文本 String message = scanner.nextLine(); // 发布消息到交换机 channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes("UTF-8")); // 输出到控制台,表示消息已发送 System.out.println(" [x] Sent '" + message + "'"); } } } }
消费者
要注意把队列和交换机绑定即可, channel.queueBind(队列名称, 交换机名称, "");
▼java复制代码import com.rabbitmq.client.*; public class ReceiveLogs { private static final String EXCHANGE_NAME = "fanout"; private static final String FANOUT_QUEUE_ONE = "logsRepo"; private static final String FANOUT_QUEUE_TWO = "mirrorRepo"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); Connection connection = factory.newConnection(); // 创建两个队列,一个专们给日志,另一个备份 Channel channel = connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); channel.queueDeclare(FANOUT_QUEUE_ONE, false, false, false, null); channel.queueBind(FANOUT_QUEUE_ONE, EXCHANGE_NAME, ""); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [x] Received 消费者1取出:"+ message + "'"); }; channel.basicConsume(FANOUT_QUEUE_ONE, true, deliverCallback, consumerTag -> { }); Channel channel2 = connection.createChannel(); channel2.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); channel2.queueDeclare(FANOUT_QUEUE_TWO, false, false, false, null); channel2.queueBind(FANOUT_QUEUE_TWO, EXCHANGE_NAME, ""); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [x] Received 消费者2取出:"+ message + "'"); }; channel.basicConsume(FANOUT_QUEUE_TWO, true, deliverCallback2, consumerTag -> { }); } }
3.4 direct交换机
特点:按照交换机上的路由表转发给指定的队列
生产者:
声明交换机和类型channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
▼java复制代码import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import java.util.Scanner; public class EmitLogDirect { private static final String EXCHANGE_NAME = "direct"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername(""); factory.setPassword(""); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { // channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); // 创建一个输入扫描器,用于读取控制台输入 Scanner scanner = new Scanner(System.in); // 使用循环,每当用户在控制台输入一行文本,就将其作为消息发送 while (scanner.hasNext()) { // 读取用户在控制台输入的下一行文本 String raw = scanner.nextLine(); // 如果用户输入的文本行包含两个单词,就将其作为消息和路由键发送 String[] s = raw.split(" "); if(s.length != 2){ continue; } String message = s[0]; String routerKey = s[1]; // 发布消息到交换机 channel.basicPublish(EXCHANGE_NAME, routerKey, null, message.getBytes("UTF-8")); // 输出到控制台,表示消息已发送 System.out.println(" [x] Sent " + message + " by direct_exchanger"); } } } }
消费者:
公式:1)交换机声明
2)队列声明
3)交换机、队列绑定,此时要声明routerKey以便于交换机定向转发
▼java复制代码import com.rabbitmq.client.*; public class ReceiveLogsDirect { private static final String EXCHANGE_NAME = "direct"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); Connection connection = factory.newConnection(); // 声明两个队列,一个员工小郑任务队列,一个员工鱼皮任务队列 Channel channel = connection.createChannel(); // 声明交换机,并指定其类型为direct channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); String queueName = "小郑的任务单"; channel.queueDeclare(queueName, false, false, false, null); // 绑定队列到交换机,指定routingKey channel.queueBind(queueName, EXCHANGE_NAME, "Xzheng"); System.out.println(" [ 小郑 ] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [ 小郑 ] Received '" + delivery.getEnvelope().getRoutingKey() + "':'" + message + "'"); }; channel.basicConsume(queueName, true, deliverCallback, consumerTag -> { }); Channel c2 = connection.createChannel(); c2.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); String qN2 = "鱼皮的任务单"; c2.queueDeclare(qN2, false, false, false, null); // 绑定队列到交换机,指定routingKey c2.queueBind(qN2, EXCHANGE_NAME, "Yupi"); System.out.println(" [ 鱼皮 ] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [ 鱼皮 ] Received '" + delivery.getEnvelope().getRoutingKey() + "':'" + message + "'"); }; c2.basicConsume(qN2, true, deliverCallback2, consumerTag -> { }); } }
