# RabbitMQ **Repository Path**: D0924/rabbit-mq ## Basic Information - **Project Name**: RabbitMQ - **Description**: 记录RabbitMQ的学习中的故事 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2022-03-25 - **Last Updated**: 2022-03-26 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # RabbitMQ一篇搞懂 ## 1.1 MQ **概述** MQ:全称Message Queue(消息队列),是在消息传递过程中保存消息的容器,多用于分布式系统直接的通信。 分布式系统通信两种方式: + 直接远程调用 + 借助第三方完成间接通信 其中发送方被称为生产者(provider),接收方被称为消费者(consumer)。 ```Mermaid graph LR A系统--->远程调用--->B系统 ``` ```Mermaid graph LR A[A系统]--->B[MQ] B--->E[B系统]; ``` ## 1.2 MQ的优劣 ### 1.2.1 优势 + 应用解构 + 异步提速 + 削峰填估 ```mermaid graph LR A[用户]---->MQ([MQ]); MQ[MQ]---->B([订单系统]); B--->C[库存系统]; B--->D[支付系统]; B--->E[物流系统]; B--->F[X系统]; ``` ```mermaid graph LR A[用户]---->MQ([MQ]); MQ[MQ]---->B([订单系统]); B--->C[库存系统]; B--->D[支付系统]; B--->E[物流系统]; B--->F[X系统]; MQ[MQ]-->db[(Database)]; ``` ```mermaid graph LR U1[用户1]--->MQ([MQ]); U2[用户2]--->MQ([MQ]); U3[用户3]--->MQ([MQ]); A1[用户4]--->MQ([MQ]); MQ[MQ]-->|每秒从MQ拉取一定量的请求进行执行|A(A系统); A[A系统]-->|同时只能处若干请求不能超|db[(Database)]; ``` ### 1.2.2劣势 + 系统可用性降低 + 系统复杂度提高 + 一致性问题 ### 1.2.3 小结 总结:何时适合使用RabbitMQ? + 生产者不需要从消费者处获得反馈。 + 容许短暂的不一致性。 + 维护成本低于低于预期收入。 ## 1.3 常见的 MQ 产品 常见的MQ产品 ## 1.4 Rabbit MQ简介 **先看基础架构架构** RabbitMQ架构 需要熟悉RabbitMQ中的一些概念 + **Broker**:接收和分发消息的应用,RabbitMQ Server就是 Message Broker + **Virtual host**:出于多租户和安全因素设计的,把 AMQP 的基本组件划分到一个虚拟的分组中,类似于网络中的 namespace 概念。当多个不同的用户使用同一个 RabbitMQ server 提供的服务时,可以划分出多个vhost,每个用户在自己的 vhost 创建 exchange/queue 等 + **Connection**: publisher/consumer 和 broker 之间的 TCP 连接,断开连接的操作只会在 client 端进行,Broker 不会断开连接,除非出现网络故障或 broker 服务出现问题 + **Channel**: 如果每一次访问 RabbitMQ 都建立一个 Connection,在消息量大的时候建立 TCP Connection的开销将是巨大的,效率也较低。Channel 是在 connection 内部建立的逻辑连接,如果应用程序支持多线程,通常每个thread创建单独的 channel 进行通讯,AMQP method 包含了channel id 帮助客户端和message broker 识别 channel,所以 channel 之间是完全隔离的。Channel 作为轻量级的 Connection 极大减少了操作系统建立 TCP connection 的开销 + **Exchange**: message 到达 broker 的第一站,根据分发规则,匹配查询表中的 routing key,分发消息到queue 中去。常用的类型有:**direct (point-to-point)**, **topic (publish-subscribe)** **andfanout (multicast)** + **Queue**:消息最终被送到这里等待 consumer 取走 + **Binding**:exchange 和 queue 之间的虚拟连接,binding 中可以包含 routing key。Binding 信息被保存到 exchange 中的查询表中,用于 message 的分发依据 ### 1.4.1 Rabbit MQ 的六种工作模式 [官方文档]: https://www.rabbitmq.com/getstarted.html "RabbitMQ的六种工作模式" ![](.\assets\mq.png) ### 1.4.2 入门 HelloWorld 简单模式 + P:生产者,也就是要发送消息的程序 + C:消费者:消息的接收者,会一直等待消息到来 + queue:消息队列,图中红色部分。可以缓存消息;生产者向其中投递消息,消费者从其中取出消息 公共依赖 ```java com.rabbitmq amqp-client 5.6.0 org.apache.maven.plugins maven-compiler-plugin 3.8.0 1.8 1.8 ``` 提供者 ```java package com.rabbitMQ.producer; 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_HelloWorld { public static void main(String[] args) throws IOException, TimeoutException { // 1. 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); // 2. 设置连接参数 factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 // 3. 创建 连接对象 Connection connection = factory.newConnection(); // 4. 创建 channel Channel channel = connection.createChannel(); // 5. 创建队列 /* queueDeclare(String queue, boolean durable, boolean exclusive, boolean autoDelete, Map arguments) 参数: 1. queue:队列名称 2. durable:是否持久化,当mq重启之后,还在 3. exclusive: * 是否独占。只能有一个消费者监听这队列 * 当Connection关闭时,是否删除队列 * 4. autoDelete:是否自动删除。当没有Consumer时,自动删除掉 5. arguments:参数。 */ channel.queueDeclare("queue_1", true, false, false, null); // 6. 发送消息 /* basicPublish(String exchange, String routingKey, BasicProperties props, byte[] body) 参数: 1. exchange:交换机名称。简单模式下交换机会使用默认的 "" 2. routingKey:路由名称 3. props:配置信息 4. body:发送消息数据 */ channel.basicPublish("", "queue_1", null, "helloWorld".getBytes()); // 7. 释放资源 channel.close(); connection.close(); } } ``` 消费者 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_HelloWorld { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("queue_1", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("queue_1", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 结果 ### 1.4.3 **Work queues** 工作队列模式  与 1.4.2 相比只是多了一个消费端,两个消费者对任务是竞争的关系 应用场景: 对于任务过重或任务较多情况使用工作队列可以提高任务处理的速度。举个例子 用户注册发送短信部署多个节点,但只需要一个节点发送成功即可。 提供者 ```java package com.rabbitMQ.producer; 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_WorkQueues { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); for (int i = 0; i < 9; i++) { channel.basicPublish("", "work_queues", null, ("Hello queues "+i).getBytes()); } channel.close(); connection.close(); } } ``` 消费者 1 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_WorkQueues_1 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("queue_1", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("work_queues", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 消费者 2 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_WorkQueues_2 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("work_queues", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 结果  ### 1.4.4 **Pub/Sub** **订阅模式**  > 在订阅者模型中多了一个 交换机(Exchange) 角色 > > P 还是 生产者,不再直接发送到队列(Queue)而是发给交换机(Exchange) > > 接着还是 Queue 接受消息 缓存消息,C1 C2 消费者 > > 交换机(Exchange) 一方面接受生产者发送的消息,另一方面知道如何处理消息,例如递交给某个队列,递交给所有队列,或将消息丢弃。消息如何处理,取决于交换机(Exchange)的类型 。常见的以下三种类型 > > + Fanout:广播,将消息交给所有绑定到交换机的队列 > + Direct:定向,把消息交给符合指定routing key 的队列 > + Topic:通配符,把消息交给符合routing pattern(路由模式) 的队列 > > 交换机(Exchange)只负责转发消息,不具备存储消息的能力,如果没有任何队列与交换机绑定,或者没有符合的路由规则队列,那么相关消息将会丢失。 提供者 ```java package com.rabbitMQ.producer; import com.rabbitmq.client.BuiltinExchangeType; 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_PubSub { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 创建交换机 /* exchangeDeclare(String exchange, BuiltinExchangeType type, boolean durable, boolean autoDelete, boolean internal, Map arguments) 参数: 1. exchange:交换机名称 2. type:交换机类型 DIRECT("direct"),:定向 FANOUT("fanout"),:扇形(广播),发送消息到每一个与之绑定队列。 TOPIC("topic"),通配符的方式 HEADERS("headers");参数匹配 3. durable:是否持久化 4. autoDelete:自动删除 5. internal:内部使用。 一般false 6. arguments:参数 */ String exchangeName = "Producer_PubSub"; // 交换机名称 String queue1Name = "Producer_PubSub_Queues_1"; // 队列名称 String queue2Name = "Producer_PubSub_Queues_2"; // 队列名称 channel.exchangeDeclare(exchangeName, BuiltinExchangeType.FANOUT, true, false, false, null); // 创建队列 channel.queueDeclare(queue1Name, true, false, false, null); channel.queueDeclare(queue2Name, true, false, false, null); // 交换机和队列的绑定 channel.queueBind(queue1Name, exchangeName, ""); channel.queueBind(queue2Name, exchangeName, ""); // 发送消息 channel.basicPublish(exchangeName, "", null, ("Message Success Mode BuiltinExchangeType.FANOUT").getBytes()); // 关闭资源 channel.close(); connection.close(); } } ``` 消费者 1 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_Consumer_PubSub_1 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("Producer_PubSub_Queues_1", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 消费者 2 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_Consumer_PubSub_2 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("Producer_PubSub_Queues_2", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 结果  ### 1.4.5 **Routing** **路由模式** ![](.\assets\mq-four.png) >P 提供者,c1,c2 消费者,X交换机 > >队列还是和交换机绑定,但是不能随意绑定了需要指定一个路由key(RoutingKey) > >消息的发送在向交换机(Exchange)发送时也需要指定消息的 RoutingKey > >当 tyep = BuiltinExchangeType.DIRECT 时 不再把消息交给每一个队列,而是根据 RoutingKey 进行判断,只有当 RoutingKey 匹配时才会发送。 > >消费者只需要与队列进行绑定,路由分发的操作已经在前面交给交换机了,所以一次消费者并不需要进行其他额外操作。 提供者 ```java package com.rabbitMQ.producer; import com.rabbitmq.client.BuiltinExchangeType; 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_Routing { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 创建交换机 /* exchangeDeclare(String exchange, BuiltinExchangeType type, boolean durable, boolean autoDelete, boolean internal, Map arguments) 参数: 1. exchange:交换机名称 2. type:交换机类型 DIRECT("direct"),:定向 FANOUT("fanout"),:扇形(广播),发送消息到每一个与之绑定队列。 TOPIC("topic"),通配符的方式 HEADERS("headers");参数匹配 3. durable:是否持久化 4. autoDelete:自动删除 5. internal:内部使用。 一般false 6. arguments:参数 */ String exchangeName = "Producer_Direct"; // 交换机名称 String queue1Name = "Producer_Direct_Queues_1"; // 队列名称 String queue2Name = "Producer_Direct_Queues_2"; // 队列名称 channel.exchangeDeclare(exchangeName, BuiltinExchangeType.DIRECT, true, false, false, null); // 创建队列 channel.queueDeclare(queue1Name, true, false, false, null); channel.queueDeclare(queue2Name, true, false, false, null); // 交换机和队列的绑定 加上路由 channel.queueBind(queue1Name, exchangeName, "error"); channel.queueBind(queue2Name, exchangeName, "error"); channel.queueBind(queue2Name, exchangeName, "info"); channel.queueBind(queue2Name, exchangeName, "warning"); // 发送消息 error 消息 channel.basicPublish(exchangeName, "error", null, ("Message Success Mode BuiltinExchangeType.DIRECT Result error").getBytes()); // 发送 info 消息 channel.basicPublish(exchangeName, "info", null, ("Message Success Mode BuiltinExchangeType.DIRECT Result info").getBytes()); // 发送 success 消息 channel.basicPublish(exchangeName, "success", null, ("Message Success Mode BuiltinExchangeType.DIRECT success").getBytes()); // 关闭资源 channel.close(); connection.close(); } } ``` 消费者 1 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_Consumer_Routing_1 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("Producer_Direct_Queues_1", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 消费者 2 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_Consumer_Routing_2 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("Producer_Direct_Queues_2", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 结果  ### 1.4.6 **Topics** **通配符模式** ![](.\assets\mq-five.png) > \# :代表没有或一个或多个单词(单词与单词之间用“.”分割) > > \*代表一个零个或一个单词; 提供者 ```java package com.rabbitMQ.producer; import com.rabbitmq.client.BuiltinExchangeType; 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_Topics { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 交换机名称 String exchangeName = "exchange_topic"; // 队列名称 String queue1Name = "topic_queue1"; String queue2Name = "topic_queue2"; // 创建交换机 channel.exchangeDeclare(exchangeName, BuiltinExchangeType.TOPIC, true, false, false, null); // 创建队列 channel.queueDeclare(queue1Name, false, false, false, null); channel.queueDeclare(queue2Name, false, false, false, null); // 绑定交换机和对垒 channel.queueBind(queue1Name,exchangeName,"*.orange.*"); channel.queueBind(queue2Name,exchangeName,"*.*.rabbite"); channel.queueBind(queue2Name,exchangeName,"Lazy.#"); // 发送消息 指定 channel.basicPublish(exchangeName, "a.orange.b", null, ("Hello queues ==a.orange.b==").getBytes()); channel.basicPublish(exchangeName, "a.b.rabbite", null, ("Hello queues ==a.b.rabbite==").getBytes()); channel.basicPublish(exchangeName, "Lazy.aabb", null, ("Hello queues ==Lazy.aabb==").getBytes()); // 未被绑定到指定通配符 消息将被丢弃 channel.basicPublish(exchangeName, "Lazy.aabb", null, ("Hello queues ==Lazy.aabb==").getBytes()); // 关闭资源 channel.close(); connection.close(); } } ``` 消费者 1 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_Consumer_Topic_1 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("topic_queue1", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 消费者 2 号 ```java package com.rabbitMQ.consume; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; /** * 接受消息的类 */ public class Consume_Consumer_Topic_2 { public static void main(String[] args) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.23.129"); // 设置连接 ip 默认 localhost factory.setPort(5672); // 设置连接端口 默认 5672 factory.setVirtualHost("/root"); // 虚拟机 我们这里使用 /root factory.setUsername("root"); // 拥有权限的用户 factory.setPassword("root"); // 拥有权限的用户密码 Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); channel.queueDeclare("work_queues", true, false, false, null); /* basicConsume(String queue, boolean autoAck, Consumer callback) 参数: 1. queue:队列名称 2. autoAck:是否自动确认 3. callback:回调对象 */ channel.basicConsume("topic_queue2", true, new DefaultConsumer(channel) { /** * @param consumerTag 标识 * @param envelope 获取一些信息,交换机,路由key... * @param properties 配置信息 * @param body 数据 * @throws IOException */ @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { System.out.println(new String(body)); } }); } } ``` 结果  ### 1.4.7 远程调用  ## 1.5 工作模式总结 简单模式 HelloWorld.java 栗子 + 一个生产者,一个消费者,不需要设置交换机使用默认的交换机 工作队列模式 Work Queue.java 栗子 + 一个生产者多个消费者(竞争关系),不需要进行设置交换机 发布订阅模式 Publish/subscribe 栗子 + 需要设置类型为 fanout 的交换机,并且交换机和队列进行绑定,当发送消息到交换机后,交换机会将消息发送到绑定的队列 路由模式 Routing 栗子 + 需要设置类型为 direct 的交换机,交换机和队列进行绑定,并且指定 routing key,当发送消息到交换机后,交换机会根据 routing key 将消息发送到对应的队列。 通配符模式 Topic 栗子 + 需要设置类型为 topic 的交换机,交换机和队列进行绑定,并且指定通配符方式的 routing key,当发送消息到交换机后,交换机会根据 routing key 将消息发送到对应的队列。 ## 1.6 spring整合RabbitMQ ### 1.6.1 提供者依赖 ```xml-dtd 4.0.0 com.fsx spring-rabbitMQ-provider 1.0-SNAPSHOT 8 8 org.springframework spring-context 5.1.7.RELEASE org.springframework.amqp spring-rabbit 2.1.8.RELEASE junit junit 4.12 org.springframework spring-test 5.1.7.RELEASE org.apache.maven.plugins maven-compiler-plugin 3.8.0 1.8 1.8 ``` ### 1.6.2 rabbitmq.properties ```java rabbitmq.host=192.168.23.129 rabbitmq.port=5672 rabbitmq.username=root rabbitmq.password=root rabbitmq.virtual-host=/root ``` ### 1.6.3 spring-rabbitmq-producer.xml ```xml-dtd ``` ### 1.6.4 ProducerTest.java ```java package com.fsx; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration(locations = "classpath:spring-rabbitmq-producer.xml") public class ProducerTest { // 1. 注入模板 @Autowired private RabbitTemplate rabbitTemplate; // 发送消息 @Test public void test() { rabbitTemplate.convertAndSend("spring_queue", "hello..."); } // Fanout 消息 @Test public void testFanout(){ rabbitTemplate.convertAndSend("spring_fanout_exchange","hello rabbitMQ 我是广播"); } // Topics 消息 @Test public void testTopics(){ rabbitTemplate.convertAndSend("spring_fanout_exchange","todo.a","spring topic.... todo.a"); rabbitTemplate.convertAndSend("spring_fanout_exchange","todo.aabbcc","spring topic.... todo.aabbcc"); rabbitTemplate.convertAndSend("spring_fanout_exchange","rabbit.a","spring topic.... rabbit.a"); } } ``` ### 1.6.5 消费者依赖 ```xml-dtd 4.0.0 com.fsx spring-rabbitMQ-consumer 1.0-SNAPSHOT 8 8 org.springframework spring-context 5.1.7.RELEASE org.springframework.amqp spring-rabbit 2.1.8.RELEASE junit junit 4.12 org.springframework spring-test 5.1.7.RELEASE org.apache.maven.plugins maven-compiler-plugin 3.8.0 1.8 1.8 ``` ### 1.6.6 spring-rabbitmq-consumer.xml ```xml-dtd ``` ### 1.6.7 SpringQueueListener.java ```java package com.fsx.rabbitmq.listener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; public class SpringQueueListener implements MessageListener { @Override public void onMessage(Message message) { //打印消息 System.out.println(new String(message.getBody())); } } ``` ### 1.6.8 ConsumerTest.java ```java package com.fsx; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration(locations = "classpath:spring-rabbitmq-consumer.xml") public class ConsumerTest { @Test public void test1(){ while (true){ } } } ``` ## 1.7 springBoot整合RabbitMQ todo.... ## 1.8 RabbitMQ 高级特性 ### 1.8.1 消息的可靠性投递 ### 1.8.2 **Consumer Ack** ### 1.8.3 消费端限流 ### 1.8.4 TTL ### 1.8.5 死信队列 ### 1.8.6 延迟队列 ## 1.9 RabbitMQ 应用集群