RabbitMQ 入门
介绍
RabbitMQ
MQ全称为Message Queue,即消息队列, RabbitMQ是由erlang语言开发,基于AMQP(Advanced Message
Queue 高级消息队列协议)协议实现的消息队列,它是一种应用程序之间的通信方法,消息队列在分布式系统开
发中应用非常广泛。RabbitMQ官方地址:http://www.rabbitmq.com/
开发中消息队列通常有如下应用场景:
1、任务异步处理。
将不需要同步处理的并且耗时长的操作由消息队列通知消息接收方进行异步处理。提高了应用程序的响应时间。
2、应用程序解耦合
MQ相当于一个中介,生产方通过MQ与消费方交互,它将应用程序进行解耦合。
市场上还有哪些消息队列?
ActiveMQ,RabbitMQ,ZeroMQ,Kafka,MetaMQ,RocketMQ、Redis。
为什么使用RabbitMQ呢?
1、使得简单,功能强大。
2、基于AMQP协议。
3、社区活跃,文档完善。
4、高并发性能好,这主要得益于Erlang语言。
5、Spring Boot默认已集成RabbitMQ
其它相关知识
AMQP是什么 ?
总结:AMQP是一套公开的消息队列协议,最早在2003年被提出,它旨在从协议层定义消息通信数据的标准格式,
为的就是解决MQ市场上协议不统一的问题。RabbitMQ就是遵循AMQP标准协议开发的MQ服务。
JMS是什么 ?
总结:
JMS是java提供的一套消息服务API标准,其目的是为所有的java应用程序提供统一的消息通信的标准,类似java的
jdbc,只要遵循jms标准的应用程序之间都可以进行消息通信。它和AMQP有什么 不同,jms是java语言专属的消
息服务标准,它是在api层定义标准,并且只能用于java应用;而AMQP是在协议层定义的标准,是跨语言的 。
入门知识
RabbitMQ的工作原理
下图是RabbitMQ的基本结构:
组成部分说明如下:
Broker:消息队列服务进程,此进程包括两个部分:Exchange和Queue。
Exchange:消息队列交换机,按一定的规则将消息路由转发到某个队列,对消息进行过虑。
Queue:消息队列,存储消息的队列,消息到达队列并转发给指定的消费方。
Producer:消息生产者,即生产方客户端,生产方客户端将消息发送到MQ。
Consumer:消息消费者,即消费方客户端,接收MQ转发的消息。
消息发布接收流程:
—–发送消息—–
1、生产者和Broker建立TCP连接。
2、生产者和Broker建立通道。
3、生产者通过通道消息发送给Broker,由Exchange将消息进行转发。
4、Exchange将消息转发到指定的Queue(队列)
—-接收消息—–
1、消费者和Broker建立TCP连接
2、消费者和Broker建立通道
3、消费者监听指定的Queue(队列)
4、当有消息到达Queue时Broker默认将消息推送给消费者。
5、消费者接收到消息。
下载安装
下载安装
RabbitMQ由Erlang语言开发,Erlang语言用于并发及分布式系统的开发,在电信领域应用广泛,OTP(Open
Telecom Platform)作为Erlang语言的一部分,包含了很多基于Erlang开发的中间件及工具库,安装RabbitMQ需
要安装Erlang/OTP,并保持版本匹配,如下图:
RabbitMQ的下载地址:http://www.rabbitmq.com/download.html
注:上图是版本对应关系
本项目使用Erlang/OTP 20.3版本和RabbitMQ3.7.3版本。
1)下载erlang
地址如下:
http://erlang.org/download/otp_win64_20.3.exe
或去老师提供的软件包中找到 otp_win64_20.3.exe,以管理员方式运行此文件,安装。
erlang安装完成需要配置erlang环境变量: ERLANG_HOME=D:\Program Files\erl9.3 在path中添
加%ERLANG_HOME%\bin;
2)安装RabbitMQ
https://github.com/rabbitmq/rabbitmq-server/releases/tag/v3.7.3
或去老师提供的软件包中找到 rabbitmq-server-3.7.3.exe,以管理员方式运行此文件,安装。
启动
安装成功后会自动创建RabbitMQ服务并且启动。
1)从开始菜单启动RabbitMQ
完成在开始菜单找到RabbitMQ的菜单:
2)如果没有开始菜单则进入安装目录下sbin目录手动启动:
3)安装并运行服务
rabbitmq-service.bat install 安装服务 rabbitmq-service.bat stop 停止服务 rabbitmq-service.bat start 启动服务
4)安装管理插件
安装rabbitMQ的管理插件,方便在浏览器端管理RabbitMQ
管理员身份运行 rabbitmq-plugins.bat enable rabbitmq_management
5) 启动成功 登录RabbitMQ
进入浏览器,输入:http://localhost:15672
初始账号和密码:guest/guest
注意事项
1、安装erlang和rabbitMQ以管理员身份运行。
2、当卸载重新安装时会出现RabbitMQ服务注册失败,此时需要进入注册表清理erlang
搜索RabbitMQ、ErlSrv,将对应的项全部删除。
Hello World案例
按照官方教程(http://www.rabbitmq.com/getstarted.html)测试hello world:
搭建环境
1)java client
生产者和消费者都属于客户端,rabbitMQ的java客户端如下:
我们先用 rabbitMQ官方提供的java client测试,目的是对RabbitMQ的交互过程有个清晰的认识。
参考 :https://github.com/rabbitmq/rabbitmq-java-client/
2)创建maven工程
创建生产者工程和消费者工程,分别加入RabbitMQ java client的依赖。
test-rabbitmq-producer:生产者工程
test-rabbitmq-consumer:消费者工程
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp‐client</artifactId>
<version>4.0.3</version><!‐‐此版本与spring boot 1.5.9版本匹配‐‐>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring‐boot‐starter‐logging</artifactId>
</dependency>
先用amqp‐client做测试,而且编写的都是官方提供的底层代码,后期整合springboot之后就更方便了
生产者
在生产者工程下的test中创建测试类如下:
package com.xuecheng.test.rabbitmq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* rabbitMQ入门程序,生产者
* @author Hr|黄锐
* @create 2019/8/8 16:11
*/
public class Producer01 {
//队列
private static final String QUEUE = "helloworld";
public static void main(String[] args) throws IOException, TimeoutException {
//建立连接(通过连接工厂)
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
//设置虚拟机,这里虚拟机类似可以同时使用多个mq,每个虚拟机相当于一个独立的mq,默认为/
//因此不必下载多个mq
connectionFactory.setVirtualHost("/");
//建立新连接
Connection connection =null;
Channel channel = null;
try {
connection = connectionFactory.newConnection();
//创建通道
channel = connection.createChannel();
//声明队列,没有就会自动创建
/*参数:
*1.队列名称
* 2.是否持久化,持久化就是重启后队列还在
* 3.是否独占连接,队列只允许在该链接中访问,如果连接关闭,队列就跟着删除了,如果设置为true,可用于临时队列的使用(测试)
* 4.自动删除,不再使用时自动删除队列,同上临队列
* 5.额外的扩展参数
*/
channel.queueDeclare(QUEUE,true,false,false,null);
//发送消息
/*参数:
*1.交换机,不指定使用默认交换机
* 2.路由key,交换机根据key将消息转发到指定的队列,默认交换机,路由key要设置为队列名称
* 3.prop,消息属性
* 4.消息的内容
*/
//消息内容
String message = "hello world hr";
channel.basicPublish("",QUEUE,null,message.getBytes());
System.out.println("send to mq "+message);
} catch (Exception e) {
e.printStackTrace();
} finally {
if (channel != null) {
channel.close();
}
if (connection != null) {
connection.close();
}
}
}
}
消费者
在消费者工程下的test中创建测试类如下:
package com.xuecheng.test.rabbitmq;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* rabbitMQ入门程序,消费者
*
* @author Hr|黄锐
* @create 2019/8/8 17:00
*/
public class Consumer {
//队列
private static final String QUEUE="helloworld";
public static void main(String[] args) throws IOException, TimeoutException {
//新建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
connectionFactory.setVirtualHost("/");
try {
//新建连接
Connection connection = connectionFactory.newConnection();
//新建通道
Channel channel = connection.createChannel();
//声明一个队列,因为生产者可能没有启动,消费者先请求,那就回会报错了,所以监听前先声明,且队列名称相同
channel.queueDeclare(QUEUE, true, false, false, null);
//消费方法
DefaultConsumer defaultConsumer=new DefaultConsumer(channel){
/**
* 当接收到消息后此方法将被调用
* @param consumerTag 消费者标签,用来标识消费者,在basicConsume可以设置
* @param envelope 信封,可以获取很多信息,交换机,通道等等
* @param properties 属性
* @param body 消息内容
* @throws IOException
*/
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
//交换机
String exchange = envelope.getExchange();
//消息id,mq在channel中用来标识消息的id,可用于编程实现确认消息已接收
long deliveryTag = envelope.getDeliveryTag();
//消息内容
String message = new String(body, "utf-8");
System.out.println("接收到:"+message);
}
};
//监听队列
/*参数:
*1.队列名称
* 2.自动回复,消费者接受到消息之后,告诉mq消息成功接收,true为自动回复mq,
* false那就要手动实现,不回复消息就会一直发
* 3.消费方法,当消费这接收到消息要执行的业务
*/
channel.basicConsume(QUEUE,true,defaultConsumer);
} catch (Exception e) {
e.printStackTrace();
}
//消费者不用关连接,因为消费者要一直监听,生产者发送完消息就可以关掉了
}
}
总结
1、发送端操作流程
1)创建连接
2)创建通道
3)声明队列
4)发送消息
2、接收端
1)创建连接
2)创建通道
3)声明队列
4)监听队列
5)接收消息
6)ack回复
注意: 1.为防止有一端先启动了,而没有声明队列等信息而导致消息传递失败,因此在生产者和消费者两端同时都要声明
RabiitMQ几大工作模式介绍
RabbitMQ有以下几种工作模式 :
1、Work queues
2、Publish/Subscribe
3、Routing
4、Topics
5、Header
6、RPC
工作模式
Work queues
work queues与入门程序相比,多了一个消费端,两个消费端共同消费同一个队列中的消息。
应用场景:对于 任务过重或任务较多情况使用工作队列可以提高任务处理的速度。
测试:
1、使用入门程序,启动多个消费者。
2、生产者发送多个消息。
结果:
1、一条消息只会被一个消费者接收;
2、rabbit采用轮询的方式将消息是平均发送给消费者的;
3、消费者在处理完某条消息后,才会收到下一条消息。
也就是在入门程序基础上可启动多个消费者线程,但是发过来的消息只能是其中一个进行处理
Publish/subscribe
工作模式
发布订阅模式:
1、每个消费者监听自己的队列。
2、生产者将消息发给broker,由交换机将消息转发到绑定此交换机的每个队列,每个绑定交换机的队列都将接收
到消息
代码
案例:
用户通知,当用户充值成功或转账完成系统通知用户,通知方式有短信、邮件多种方法 。
1、生产者
声明Exchange_fanout_inform交换机。
声明两个队列并且绑定到此交换机,绑定时不需要指定routingkey
发送消息时不需要指定routingkey
package com.xuecheng.test.rabbitmq;
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;
/**
* 发布订阅模式
* 可实现多个消费者接受同一个消息
* @author Hr|黄锐
* @create 2019/8/8 20:22
*/
public class Producer_publish {
//队列
private static final String QUEUE_EMAIL = "email";
private static final String QUEUE_SMS = "sms";
//交换机
private static final String EXCHANGE_FANOUT_INFORM = "exchange_fanout_inform";
public static void main(String[] args) throws IOException, TimeoutException {
//建立连接(通过连接工厂)
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
//设置虚拟机,这里虚拟机类似可以同时使用多个mq,每个虚拟机相当于一个独立的mq,默认为/
//因此不必下载多个mq
connectionFactory.setVirtualHost("/");
//建立新连接
Connection connection =null;
Channel channel = null;
try {
connection = connectionFactory.newConnection();
//创建通道
channel = connection.createChannel();
//声明队列,没有就会自动创建
channel.queueDeclare(QUEUE_EMAIL,true,false,false,null);
channel.queueDeclare(QUEUE_SMS,true,false,false,null);
/**
* 参数明细:
* 交换机类型:
* fanout就是对应发布订阅模式
*/
channel.exchangeDeclare(EXCHANGE_FANOUT_INFORM, BuiltinExchangeType.FANOUT);
//队列绑定交换机
/**参数明细:
* 1.queue:队列名称
* 2.exchange:交换机名称
* 3.路由key:在发布订阅模式定为""空串
*/
channel.queueBind(QUEUE_EMAIL,EXCHANGE_FANOUT_INFORM,"");
channel.queueBind(QUEUE_SMS,EXCHANGE_FANOUT_INFORM,"");
//发送消息
for (int i = 0; i < 5; i++) {
//消息内容
String message = "send inform message to user";
channel.basicPublish(EXCHANGE_FANOUT_INFORM,"",null,message.getBytes());
System.out.println("send to mq "+message);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
if (channel != null) {
channe.close();
}
if (connection != null) {
connection.close();
}
}
}
}
2、邮件发送消费者
package com.xuecheng.test.rabbitmq;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* 发布订阅模式,订阅方
* @author Hr|黄锐
* @create 2019/8/8 17:00
*/
public class Consumer_subscribe_email {
//队列
private static final String QUEUE_EMAIL = "email";
private static final String QUEUE_SMS = "sms";
//交换机
private static final String EXCHANGE_FANOUT_INFORM = "exchange_fanout_inform";
public static void main(String[] args) throws IOException, TimeoutException {
//新建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
connectionFactory.setVirtualHost("/");
try {
//新建连接
Connection connection = connectionFactory.newConnection();
//新建通道
Channel channel = connection.createChannel();
//声明一个队列,因为生产者可能没有启动,消费者先请求,那就回会报错了,所以监听前先声明,且队列名称相同
channel.queueDeclare(QUEUE_EMAIL, true, false, false, null);
//声明交换机
channel.exchangeDeclare(EXCHANGE_FANOUT_INFORM, BuiltinExchangeType.FANOUT);
//绑定队列
channel.queueBind(QUEUE_EMAIL, EXCHANGE_FANOUT_INFORM, "");
//消费方法
DefaultConsumer defaultConsumer=new DefaultConsumer(channel){
/**
* 当接收到消息后此方法将被调用
* @param consumerTag 消费者标签,用来标识消费者,在basicConsume可以设置
* @param envelope 信封,可以获取很多信息,交换机,通道等等
* @param properties 属性
* @param body 消息内容
* @throws IOException
*/
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
//交换机
String exchange = envelope.getExchange();
//消息id,mq在channel中用来标识消息的id,可用于编程实现确认消息已接收
long deliveryTag = envelope.getDeliveryTag();
//消息内容
String message = new String(body, "utf-8");
System.out.println("接收到:"+message);
}
};
//监听队列
/*参数:
*1.队列名称
* 2.自动回复,消费者接受到消息之后,告诉mq消息成功接收,true为自动回复mq,
* false那就要手动实现,不回复消息就会一直发
* 3.消费方法,当消费这接收到消息要执行的业务
*/
channel.basicConsume(QUEUE_EMAIL,true,defaultConsumer);
} catch (Exception e) {
e.printStackTrace();
}
//消费者不用关连接,因为消费者要一直监听,生产者发送完消息就可以关掉了
}
}
按照上边的代码,编写邮件通知的消费代码。
3、短信发送消费者
参考上边的邮件发送消费者代码编写。
和工作模式的区别就是,每个消费者管理一个自己的队列,生产者发送消息时都能接收到,同时也可改装为work工作模式,加上多个消费者
思考
1、publish/subscribe与work queues有什么区别。
区别:
1)work queues不用定义交换机,而publish/subscribe需要定义交换机,且类型为fanout。
2)publish/subscribe的生产方是面向交换机发送消息,work queues的生产方是面向队列发送消息(底层使用默认
交换机)。
3)publish/subscribe需要设置队列和交换机的绑定,work queues不需要设置,实质上work queues会将队列绑
定到默认的交换机 。
相同点:
所以两者实现的发布/订阅的效果是一样的,多个消费端监听同一个队列不会重复消费消息。
2、实质工作用什么 publish/subscribe还是work queues。
建议使用 publish/subscribe,发布订阅模式比工作队列模式更强大,并且发布订阅模式可以指定自己专用的交换
机。
Routing
工作模式
路由模式:
1、每个消费者监听自己的队列,并且设置routingkey。
2、生产者将消息发给交换机,由交换机根据routingkey来转发消息到指定的队列。
代码
1、生产者
声明exchange_routing_inform交换机。
声明两个队列并且绑定到此交换机,绑定时需要指定routingkey
发送消息时需要指定routingkey
package com.xuecheng.test.rabbitmq;
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;
/**
* 路由模式,生产者
* @author Hr|黄锐
* @create 2019/8/8 21:04
*/
public class Producer_routing {
//队列
private static final String QUEUE_EMAIL = "email";
private static final String QUEUE_SMS = "sms";
//交换机
private static final String EXCHANGE_ROUTING_INFORM = "exchange_routing_inform";
//routingKey
private static final String ROUTINGKEY_EMAIL = "inform_email";
private static final String ROUTINGKEY_SMS = "inform_sms";
public static void main(String[] args) throws IOException, TimeoutException {
//建立连接(通过连接工厂)
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
//设置虚拟机,这里虚拟机类似可以同时使用多个mq,每个虚拟机相当于一个独立的mq,默认为/
//因此不必下载多个mq
connectionFactory.setVirtualHost("/");
//建立新连接
Connection connection =null;
Channel channel = null;
try {
connection = connectionFactory.newConnection();
//创建通道
channel = connection.createChannel();
//声明队列,没有就会自动创建
channel.queueDeclare(QUEUE_EMAIL,true,false,false,null);
channel.queueDeclare(QUEUE_SMS,true,false,false,null);
/**
* 参数明细:
* 交换机类型:
* fanout就是对应发布订阅模式
*/
channel.exchangeDeclare(EXCHANGE_ROUTING_INFORM, BuiltinExchangeType.DIRECT);
//队列绑定交换机
/**参数明细:
* 1.queue:队列名称
* 2.exchange:交换机名称
* 3.路由key:在发布订阅模式定为""空串
*/
channel.queueBind(QUEUE_EMAIL,EXCHANGE_ROUTING_INFORM,ROUTINGKEY_EMAIL);
channel.queueBind(QUEUE_SMS,EXCHANGE_ROUTING_INFORM,ROUTINGKEY_SMS);
//发送消息
for (int i = 0; i < 5; i++) {
//消息内容
//发送消息时要指定routingkey
String message = "send inform message to user";
channel.basicPublish(EXCHANGE_ROUTING_INFORM,ROUTINGKEY_EMAIL,null,message.getBytes());
System.out.println("send to mq "+message);
}
for (int i = 0; i < 5; i++) {
//消息内容
//发送消息时要指定routingkey
String message = "send inform message to user";
channel.basicPublish(EXCHANGE_ROUTING_INFORM,ROUTINGKEY_SMS,null,message.getBytes());
System.out.println("send to mq "+message);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
if (channel != null) {
channel.close();
}
if (connection != null) {
connection.close();
}
}
}
}
2、邮件发送消费者
package com.xuecheng.test.rabbitmq;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* 路由模式,消费方
* @author Hr|黄锐
* @create 2019/8/8 17:00
*/
public class Consumer_routing_email {
//队列
private static final String QUEUE_EMAIL = "email";
//交换机
private static final String EXCHANGE_ROUTING_INFORM = "exchange_routing_inform";
//routingKey
private static final String ROUTINGKEY_EMAIL = "inform_email";
public static void main(String[] args) throws IOException, TimeoutException {
//新建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
connectionFactory.setVirtualHost("/");
try {
//新建连接
Connection connection = connectionFactory.newConnection();
//新建通道
Channel channel = connection.createChannel();
//声明一个队列,因为生产者可能没有启动,消费者先请求,那就回会报错了,所以监听前先声明,且队列名称相同
channel.queueDeclare(QUEUE_EMAIL, true, false, false, null);
//声明交换机
channel.exchangeDeclare(EXCHANGE_ROUTING_INFORM, BuiltinExchangeType.DIRECT);
//绑定队列
channel.queueBind(QUEUE_EMAIL, EXCHANGE_ROUTING_INFORM, ROUTINGKEY_EMAIL);
//消费方法
DefaultConsumer defaultConsumer=new DefaultConsumer(channel){
/**
* 当接收到消息后此方法将被调用
* @param consumerTag 消费者标签,用来标识消费者,在basicConsume可以设置
* @param envelope 信封,可以获取很多信息,交换机,通道等等
* @param properties 属性
* @param body 消息内容
* @throws IOException
*/
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
//交换机
String exchange = envelope.getExchange();
//消息id,mq在channel中用来标识消息的id,可用于编程实现确认消息已接收
long deliveryTag = envelope.getDeliveryTag();
//消息内容
String message = new String(body, "utf-8");
System.out.println("接收到:"+message);
}
};
//监听队列
/*参数:
*1.队列名称
* 2.自动回复,消费者接受到消息之后,告诉mq消息成功接收,true为自动回复mq,
* false那就要手动实现,不回复消息就会一直发
* 3.消费方法,当消费这接收到消息要执行的业务
*/
channel.basicConsume(QUEUE_EMAIL,true,defaultConsumer);
} catch (Exception e) {
e.printStackTrace();
}
//消费者不用关连接,因为消费者要一直监听,生产者发送完消息就可以关掉了
}
}
3、短信发送消费者 参考邮件发送消费者的代码流程,编写短信通知的代码。
测试
打开RabbitMQ的管理界面,观察交换机绑定情况:
使用生产者发送若干条消息,交换机根据routingkey转发消息到指定的队列。
思考
1、Routing模式和Publish/subscibe有啥区别? Routing模式要求队列在绑定交换机时要指定routingkey,消息会转发到符合routingkey的队列。
routing模式就是在发布订阅模式基础上,加了一层条件判断的routingkey,符合条件才能通过
Topics
工作模式
路由模式:
1、每个消费者监听自己的队列,并且设置带统配符的routingkey。
2、生产者将消息发给broker,由交换机根据routingkey来转发消息到指定的队列。
代码
案例: 根据用户的通知设置去通知用户,设置接收Email的用户只接收Email,设置接收sms的用户只接收sms,设置两种 通知类型都接收的则两种通知都有效。
1、生产者 声明交换机,指定topic类型:
/**
* 声明交换机
* param1:交换机名称
* param2:交换机类型 四种交换机类型:direct、fanout、topic、headers
*/
channel.exchangeDeclare(EXCHANGE_TOPICS_INFORM, BuiltinExchangeType.TOPIC);
//Email通知
channel.basicPublish(EXCHANGE_TOPICS_INFORM, "inform.email", null, message.getBytes());
//sms通知
channel.basicPublish(EXCHANGE_TOPICS_INFORM, "inform.sms", null, message.getBytes());
//两种都通知
channel.basicPublish(EXCHANGE_TOPICS_INFORM, "inform.sms.email", null, message.getBytes());
完整代码:
package com.xuecheng.test.rabbitmq;
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;
/**
* topics模式,通配符匹配
* @author Hr|黄锐
* @create 2019/8/8 21:04
*/
public class Producer_topics {
//队列
private static final String QUEUE_EMAIL = "email";
private static final String QUEUE_SMS = "sms";
//交换机
private static final String EXCHANGE_TOPICS_INFORM = "exchange_topics_inform";
//routingKey
private static final String ROUTINGKEY_EMAIL = "inform.#.email.#";//inform.email
private static final String ROUTINGKEY_SMS = "inform.#.sms.#";
public static void main(String[] args) throws IOException, TimeoutException {
//建立连接(通过连接工厂)
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
//设置虚拟机,这里虚拟机类似可以同时使用多个mq,每个虚拟机相当于一个独立的mq,默认为/
//因此不必下载多个mq
connectionFactory.setVirtualHost("/");
//建立新连接
Connection connection =null;
Channel channel = null;
try {
connection = connectionFactory.newConnection();
//创建通道
channel = connection.createChannel();
//声明队列,没有就会自动创建
channel.queueDeclare(QUEUE_EMAIL,true,false,false,null);
channel.queueDeclare(QUEUE_SMS,true,false,false,null);
/**
* 参数明细:
* 交换机类型:
* fanout就是对应发布订阅模式
*/
channel.exchangeDeclare(EXCHANGE_TOPICS_INFORM, BuiltinExchangeType.TOPIC);
//队列绑定交换机
/**参数明细:
* 1.queue:队列名称
* 2.exchange:交换机名称
* 3.路由key:在发布订阅模式定为""空串
*/
channel.queueBind(QUEUE_EMAIL,EXCHANGE_TOPICS_INFORM,ROUTINGKEY_EMAIL);
channel.queueBind(QUEUE_SMS,EXCHANGE_TOPICS_INFORM,ROUTINGKEY_SMS);
//发送消息
for (int i = 0; i < 5; i++) {
//消息内容
//发送消息时要指定routingkey
String message = "send inform message to user";
channel.basicPublish(EXCHANGE_TOPICS_INFORM,"inform.email",null,message.getBytes());
System.out.println("send to mq "+message);
}
for (int i = 0; i < 5; i++) {
//消息内容
//发送消息时要指定routingkey
String message = "send inform message to user";
channel.basicPublish(EXCHANGE_TOPICS_INFORM,"inform.sms",null,message.getBytes());
System.out.println("send to mq "+message);
}
for (int i = 0; i < 5; i++) {
//消息内容
//发送消息时要指定routingkey
String message = "send inform message to user";
channel.basicPublish(EXCHANGE_TOPICS_INFORM,"inform.sms.email",null,message.getBytes());
System.out.println("send to mq "+message);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
if (channel != null) {
channel.close();
}
if (connection != null) {
connection.close();
}
}
}
}
2、消费端
队列绑定交换机指定通配符:
统配符规则:
中间以“.”分隔。
符号#可以匹配多个词,符号*可以匹配一个词语。
//声明队列
channel.queueDeclare(QUEUE_INFORM_EMAIL, true, false, false, null);
channel.queueDeclare(QUEUE_INFORM_SMS, true, false, false, null);
//声明交换机
channel.exchangeDeclare(EXCHANGE_TOPICS_INFORM, BuiltinExchangeType.TOPIC);
//绑定email通知队列
channel.queueBind(QUEUE_INFORM_EMAIL,EXCHANGE_TOPICS_INFORM,"inform.#.email.#");
//绑定sms通知队列
channel.queueBind(QUEUE_INFORM_SMS,EXCHANGE_TOPICS_INFORM,"inform.#.sms.#");
测试
使用生产者发送若干条消息,交换机根据routingkey统配符匹配并转发消息到指定的队列。
思考
1、本案例的需求使用Routing工作模式能否实现?
使用Routing模式也可以实现本案例,共设置三个 routingkey,分别是email、sms、all,email队列绑定email和
all,sms队列绑定sms和all,这样就可以实现上边案例的功能,实现过程比topics复杂。
Topic模式更多加强大,它可以实现Routing、publish/subscirbe模式的功能。
topics模式也就是在routing模式基础上,更灵活的判断条件,例如: inform.#.email.# 可以匹配inform.email 也可匹配 inform.aa.email 还可匹配 inform.email.ddd
Header模式
header模式与routing不同的地方在于,header模式取消routingkey,使用header中的 key/value(键值对)匹配
队列。
案例:
根据用户的通知设置去通知用户,设置接收Email的用户只接收Email,设置接收sms的用户只接收sms,设置两种
通知类型都接收的则两种通知都有效。
代码:
1)生产者
队列与交换机绑定的代码与之前不同,如下:
Map<String, Object> headers_email = new Hashtable<String, Object>();
headers_email.put("inform_type", "email");
Map<String, Object> headers_sms = new Hashtable<String, Object>();
headers_sms.put("inform_type", "sms");
channel.queueBind(QUEUE_INFORM_EMAIL,EXCHANGE_HEADERS_INFORM,"",headers_email);
channel.queueBind(QUEUE_INFORM_SMS,EXCHANGE_HEADERS_INFORM,"",headers_sms);
通知:
String message = "email inform to user"+i;
Map<String,Object> headers = new Hashtable<String, Object>();
headers.put("inform_type", "email");//匹配email通知消费者绑定的header
//headers.put("inform_type", "sms");//匹配sms通知消费者绑定的header
AMQP.BasicProperties.Builder properties = new AMQP.BasicProperties.Builder();
properties.headers(headers);
//Email通知
channel.basicPublish(EXCHANGE_HEADERS_INFORM, "", properties.build(), message.getBytes());
2)发送邮件消费者
channel.exchangeDeclare(EXCHANGE_HEADERS_INFORM, BuiltinExchangeType.HEADERS);
Map<String, Object> headers_email = new Hashtable<String, Object>();
headers_email.put("inform_email", "email");
//交换机和队列绑定
channel.queueBind(QUEUE_INFORM_EMAIL,EXCHANGE_HEADERS_INFORM,"",headers_email);
//指定消费队列
channel.basicConsume(QUEUE_INFORM_EMAIL, true, consumer);
3)测试
RPC
RPC即客户端远程调用服务端的方法 ,使用MQ可以实现RPC的异步调用,基于Direct交换机实现,流程如下:
1、客户端即是生产者就是消费者,向RPC请求队列发送RPC调用消息,同时监听RPC响应队列。
2、服务端监听RPC请求队列的消息,收到消息后执行服务端的方法,得到方法返回的结果
3、服务端将RPC方法 的结果发送到RPC响应队列
4、客户端(RPC调用方)监听RPC响应队列,接收到RPC调用结果。
Spring整合RibbitMQ
搭建SpringBoot环境
我们选择基于Spring-Rabbit去操作RabbitMQ
https://github.com/spring-projects/spring-amqp
使用spring-boot-starter-amqp会自动添加spring-rabbit依赖,如下:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring‐boot‐starter‐amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring‐boot‐starter‐test</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring‐boot‐starter‐logging</artifactId>
</dependency>