
direct交换机是RabbitMQ默认交换机。docker start rabbitmq
docker exec -it rabbitmq rabbitmq-plugins enable rabbitmq_management

<dependencies>
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.14.0</version>
</dependency>
</dependencies>
package com.knife.demo01.simple;
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 void main(String[] args) throws IOException, TimeoutException {
// 1.创建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("192.168.70.130"); // 虚拟机ip
connectionFactory.setPort(5672);
connectionFactory.setUsername("admin");
connectionFactory.setPassword("admin");
connectionFactory.setVirtualHost("/");
// 2.创建连接
Connection connection = connectionFactory.newConnection();
// 3.建立信道
Channel channel = connection.createChannel();
// 4.创建队列,如果队列已存在,则使用该队列
/**
* 参数1:队列名
* 参数2:是否持久化,true表示MQ重启后队列还在。
* 参数3:是否私有化,false表示所有消费者都可以访问,true表示只有第一次拥有它的消费者才能访问
* 参数4:是否自动删除,true表示不再使用队列时自动删除队列
* 参数5:其他额外参数
*/
channel.queueDeclare("simple_queue",false,false,false,null);
// 5.发送消息
String message = "hello!rabbitmq!";
/**
* 参数1:交换机名,""表示默认交换机direct
* 参数2:路由键,简单模式就是队列名
* 参数3:其他额外参数
* 参数4:要传递的消息字节数组
*/
channel.basicPublish("","simple_queue",null,message.getBytes());
// 6.关闭信道和连接
channel.close();
connection.close();
System.out.println("===发送成功===");
}
}
运行生产者代码结果:


package com.knife.demo01.simple;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
// 消费者
public class Consumer {
public static void main(String[] args) throws IOException, TimeoutException {
// 1.创建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("192.168.70.130"); // 虚拟机ip
connectionFactory.setPort(5672); // 端口号
connectionFactory.setUsername("admin");
connectionFactory.setPassword("admin");
connectionFactory.setVirtualHost("/");
// 2.创建连接
Connection connection = connectionFactory.newConnection();
// 3.建立信道
Channel channel = connection.createChannel();
// 4.监听队列
/**
* 参数1:监听的队列名
* 参数2:是否自动签收,如果设置为false,则需要手动确认消息已收到,否则MQ会一直发送消息
* 参数3:Consumer的实现类,重写该类方法表示接受到消息后如何消费
*/
channel.basicConsume("simple_queue",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
// 接收消息
String message = new String(body, "UTF-8");
System.out.println("接受消息,消息为:"+message);
}
});
}
}
运行消费者代码结果:


工作队列模式(Work Queue)与简单模式相比,多了一些消费者,该模式也使用direct交换机(默认使用的交换机),应用于处理消息较多的情况。特点如下:
Work Queues 对于任务过重或任务较多情况使用工作队列可以提高任务处理的速度。例如:短信服务部署多个,只需要有一个节点成功发送即可


public class Producer {
public static void main(String[] args) throws IOException, TimeoutException {
// 创建工厂连接
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130"); //虚拟机地址
cf.setPort(5672); //rabbitmq的端口号
cf.setUsername("admin"); // 用户名
cf.setPassword("admin"); // 密码
cf.setVirtualHost("/");
// 创建连接
Connection connection = cf.newConnection();
// 创建信道
Channel channel = connection.createChannel();
// 创建队列
/**
* 参数1:队列名称
* 参数二:是否持久化
* 参数三:是否私有化
* 参数4:队列使用完毕后是否自动删除
*/
channel.queueDeclare("work_queue",true,false,false,null);
// 发送消息
/**
* 参数1:交换机名,""表示默认交换机direct
* 参数2:路由键,简单模式就是队列名
* 参数3:其他额外参数
* 参数4:要传递的消息字节数组
*/
for (int i = 1; i <= 10; i++) {
channel.basicPublish("","work_queue", MessageProperties.PERSISTENT_TEXT_PLAIN,
("你好,这是今天的第"+i+"条消息").getBytes());
}
// 关闭资源
channel.close();
connection.close();
System.out.println("消息发送成功====");
}
}
public class Consumer {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
// 4. 接收消息
channel.basicConsume("work_queue",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("消费者1 = " + msg);
}
});
channel.close();
conn.close();
}
}
代码运行的效果和上面简单模式的类似。

需要不同的消费者进行不同的处理
Exchange角色,而且过程略有变化:
不再发送到队列中,而是发给X(交换机)绑定到交换机的队列符合指定routing key 的队列符合routing pattern(路由模式) 的队列,Exchange(交换机)只负责转发消息,不具备存储消息的能力,因此如果没有任何队列与 Exchange 绑定,或者没有符合路由规则的队列,那么消息会丢失!将消息转发到绑定此交换机的每个队列中(即可以转发到多个队列中)。多个队列。发布订阅模式使用fanout交换机。
public class Producer {
// 生产者
public static void main(String[] args) throws IOException, TimeoutException {
// 1.创建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("192.168.70.130");
connectionFactory.setPort(5672);
connectionFactory.setUsername("admin");
connectionFactory.setPassword("admin");
connectionFactory.setVirtualHost("/");
// 2.创建连接
Connection connection = connectionFactory.newConnection();
// 3.建立信道
Channel channel = connection.createChannel();
// 4.创建交换机
/**
* 参数1:交换机名
* 参数2:交换机类型
* 参数3:交换机持久化
*/
// BuiltinExchangeType.FANOUT:表示广播模式
channel.exchangeDeclare("exchange_fanout", BuiltinExchangeType.FANOUT, true);
// 5.创建队列
// 队列1
channel.queueDeclare("SEND_MAIL", true, false, false, null);
// 队列2
channel.queueDeclare("SEND_MESSAGE", true, false, false, null);
// 队列3
channel.queueDeclare("SEND_STATION", true, false, false, null);
// 6.交换机绑定队列
/**
* 参数1:队列的名称
* 参数2:交换机名
* 参数3:路由关键字,发布订阅模式写""即可
*/
channel.queueBind("SEND_MAIL", "exchange_fanout", "");
channel.queueBind("SEND_MESSAGE", "exchange_fanout", "");
channel.queueBind("SEND_STATION", "exchange_fanout", "");
// 7.发送消息
channel.basicPublish("exchange_fanout", "", null,
("618商品开抢了!").getBytes(StandardCharsets.UTF_8));
// for (int i = 1; i <= 10; i++) {
// channel.basicPublish("exchange_fanout", "", null,
// ("你好,尊敬的用户,秒杀商品开抢了!" + i).getBytes(StandardCharsets.UTF_8));
// }
// 8.关闭资源
channel.close();
connection.close();
}
}
public class ConsumerMail {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_MAIL",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送邮件消息 = " + msg);
}
});
}
}
public class ConsumerMessage {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_MESSAGE",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送短信消息 = " + msg);
}
});
}
}
public class ConsumerStation {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_STATION",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送站內信 = " + msg);
}
});
}
}
而一些小的促销活动为了节约成本,只发布到站内信队列。此时需要使用路由模式(Routing)完成这一需求。
要指定一个 RoutingKey(路由key),消息的发送方在向Exchange发送消息时,也必须指定消息的 RoutingKey,Exchange不再把消息交给每一个绑定的队列,而是根据消息的Routing Key进行判断,只有队列的Routingkey 与消息的 Routing key 完全一致,才会接收到消息。
其实代码还是差不多的,只是使用到的交换机不一样了,多了一个route Key。
public class Producer {
// 生产者
public static void main(String[] args) throws IOException, TimeoutException {
// 1.创建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("192.168.70.130");
connectionFactory.setPort(5672);
connectionFactory.setUsername("admin");
connectionFactory.setPassword("admin");
connectionFactory.setVirtualHost("/");
// 2.创建连接
Connection connection = connectionFactory.newConnection();
// 3.建立信道
Channel channel = connection.createChannel();
// 4.创建交换机
/**
* 参数1:交换机名
* 参数2:交换机类型
* 参数3:交换机持久化
*/
// BuiltinExchangeType.DIRECT
channel.exchangeDeclare("routing_exchange", BuiltinExchangeType.DIRECT, true);
// 5.创建队列
// 队列1
channel.queueDeclare("SEND_MAILRoute", true, false, false, null);
// 队列2
channel.queueDeclare("SEND_MESSAGERoute", true, false, false, null);
// 队列3
channel.queueDeclare("SEND_STATIONRoute", true, false, false, null);
// 6.交换机绑定队列
/**
* 参数1:队列的名称
* 参数2:交换机名
* 参数3:路由关键字(路由名称)
*/
channel.queueBind("SEND_MAILRoute", "routing_exchange", "import");
channel.queueBind("SEND_MESSAGERoute", "routing_exchange", "import");
channel.queueBind("SEND_STATIONRoute", "routing_exchange", "import");
channel.queueBind("SEND_STATIONRoute", "routing_exchange", "normal");
// 7.发送消息
channel.basicPublish("routing_exchange", "import", null,
("618商品开抢了!").getBytes(StandardCharsets.UTF_8));
channel.basicPublish("routing_exchange", "normal", null,
("normal路由的息").getBytes(StandardCharsets.UTF_8));
// for (int i = 1; i <= 10; i++) {
// channel.basicPublish("exchange_fanout", "", null,
// ("你好,尊敬的用户,秒杀商品开抢了!" + i).getBytes(StandardCharsets.UTF_8));
// }
// 8.关闭资源
channel.close();
connection.close();
}
}
public class ConsumerMail {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_MAILRoute",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送邮件消息 = " + msg);
}
});
}
}
public class ConsumerMessage {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_MESSAGERoute",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送短信消息 = " + msg);
}
});
}
}
public class ConsumerStation {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.70.130");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_STATIONRoute",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送站內信 = " + msg);
}
});
}
}

给队列绑定带通配符的路由关键字,只要消息的RoutingKey能实现通配符匹配,就会将消息转发到该队列。通配符模式比路由模式更灵活,使用topic交换机。#可以匹配任意多个单词,*可以匹配任意一个单词。

public class Producer {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工厂
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.126.10");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.创建交换机
/**
* 参数1:交换机名
* 参数2:交换机类型
* 参数3:交换机持久化
*/
channel.exchangeDeclare("exchange_topic", BuiltinExchangeType.TOPIC, true);
/**
* 参数1:队列名
* 参数2:是否持久化,true表示MQ重启后队列还在。
* 参数3:是否私有化,false表示所有消费者都可以访问,true表示只有第一次拥有它的消费者才能访问
* 参数4:是否自动删除,true表示不再使用队列时自动删除队列
* 参数5:其他额外参数
*/
//5.创建队列(举例订单给手机发送,邮箱,站内)发送消息,3个队列
channel.queueDeclare("SEND_MAIL3", true, false, false, null);
channel.queueDeclare("SEND_MESSAGE3", true, false, false, null);
channel.queueDeclare("SEND_STATION3", true, false, false, null);
//6.将队列和交换机绑定
/**
* 参数1:队列名
* 参数2:交换机名
* 参数3:路由关键字,发布订阅模式写""即可
*/
channel.queueBind("SEND_MAIL3", "exchange_topic", "#.mail.#");
channel.queueBind("SEND_MESSAGE3", "exchange_topic", "#.message.#");
channel.queueBind("SEND_STATION3", "exchange_topic", "#.station.#");
//8.发送消息
channel.basicPublish("exchange_topic", "mail.message.station", null, "618大促销活动".getBytes(StandardCharsets.UTF_8));
channel.basicPublish("exchange_topic", "station", null, "618小促销活动".getBytes(StandardCharsets.UTF_8));
//9.关闭资源
channel.close();
conn.close();
System.out.println("发送消息成功");
}
}
public class ConsumerMail {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.126.10");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_MAIL3",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送邮件消息 = " + msg);
}
});
}
}
public class ConsumerMessage {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.126.10");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_MESSAGE3",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送短信消息 = " + msg);
}
});
}
}
public class ConsumerStation {
public static void main(String[] args) throws IOException, TimeoutException {
//1.创建连接工程
ConnectionFactory cf = new ConnectionFactory();
cf.setHost("192.168.126.10");
cf.setPort(5672);
cf.setUsername("admin");
cf.setPassword("admin");
cf.setVirtualHost("/");
//2.创建连接
Connection conn = cf.newConnection();
//3.创建信道
Channel channel = conn.createChannel();
//4.监听队列
channel.basicConsume("SEND_STATION3",true,new DefaultConsumer(channel){
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
String msg = new String(body, StandardCharsets.UTF_8);
System.out.println("发送站內信 = " + msg);
}
});
}
}
RabbitMQ一、RabbitMQ的介绍与安装(docker)
RabbitMQ三、springboot整合rabbitmq(消息可靠性、高级特性)