RabbitMQ快速入门详解
原文链接:https://blog.csdn.net/kavito/article/details/91403659
在介绍RabbitMQ之前,我们先来看下面一个电商项目的场景:
-
商品的原始数据保存在数据库中,增删改查都在数据库中完成。
-
搜索服务数据来源是索引库(Elasticsearch),如果数据库商品发生变化,索引库数据不能及时更新。
-
商品详情做了页面静态化处理,静态页面数据也不会随着数据库商品更新而变化。
如果我们在后台修改了商品的价格,搜索页面和商品详情页显示的依然是旧的价格,这样显然不对。该如何解决?
我们可能会想到这么做:
-
方案1:每当后台对商品做增删改操作,同时修改索引库数据及更新静态页面。
-
方案2:搜索服务和商品页面静态化服务对外提供操作接口,后台在商品增删改后,调用接口。
这两种方案都有个严重的问题:就是代码耦合,后台服务中需要嵌入搜索和商品页面服务,违背了微服务的独立原则。
这时,我们就会采用另外一种解决办法,那就是消息队列!
商品服务对商品增删改以后,无需去操作索引库和静态页面,只需向MQ发送一条消息(比如包含商品id的消息),也不关心消息被谁接收。 搜索服务和静态页面服务监听MQ,接收消息,然后分别去处理索引库和静态页面(根据商品id去更新索引库和商品详情静态页面)。
MySQL,直接导致无数的行锁表锁,甚至最后请求会堆积过多,从而触发too many connections错误。通过使用消息队列,我们可以异步处理请求,从而缓解系统的压力。将不需要同步处理的并且耗时长的操作由消息队列通知消息接收方进行异步处理。减少了应用程序的响应时间。
2、应用程序解耦合:
MQ相当于一个中介,生产方通过MQ与消费方交互,它将应用程序进行解耦合。
http://www.rabbitmq.com/download.html
六种消息模型
127.0.0.1:15672,默认用户及密码:guest guest)
127.0.0.1:15672,默认用户及密码:guest guest)

点击队列名称,进入详情页,可以查看消息:

https://blog.csdn.net/zhu_tianwei/article/details/40923131
Exchange(交换机)只负责转发消息,不具备存储消息的能力,因此如果没有任何队列与Exchange绑定,或者没有符合路由规则的队列,那么消息会丢失!
③Publish/subscribe(交换机类型:Fanout,也称为广播 )
Publish/subscribe模型示意图 :

生产者
和前面两种模式不同:
-
1) 声明Exchange,不再声明Queue
-
2) 发送消息到Exchange,不再发送到Queue
- public class Send {
- private final static String EXCHANGE_NAME = "test_fanout_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明exchange,指定类型为fanout
- channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
- // 消息内容
- String message = "注册成功!!";
- // 发布消息到Exchange
- channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes());
- System.out.println(" [生产者] Sent '" + message + "'");
- channel.close();
- connection.close();
- }
- }
消费者1 (注册成功发给短信服务)
- public class Recv {
- private final static String QUEUE_NAME = "fanout_exchange_queue_sms";//短信队列
- private final static String EXCHANGE_NAME = "test_fanout_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明队列
- channel.queueDeclare(QUEUE_NAME, false, false, false, null);
- // 绑定队列到交换机
- channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "");
- // 定义队列的消费者
- DefaultConsumer consumer = new DefaultConsumer(channel) {
- // 获取消息,并且处理,这个方法类似事件监听,如果有消息的时候,会被自动调用
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
- byte[] body) throws IOException {
- // body 即消息体
- String msg = new String(body);
- System.out.println(" [短信服务] received : " + msg + "!");
- }
- };
- // 监听队列,自动返回完成
- channel.basicConsume(QUEUE_NAME, true, consumer);
- }
- }
消费者2(注册成功发给邮件服务)
- public class Recv2 {
- private final static String QUEUE_NAME = "fanout_exchange_queue_email";//邮件队列
- private final static String EXCHANGE_NAME = "test_fanout_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明队列
- channel.queueDeclare(QUEUE_NAME, false, false, false, null);
- // 绑定队列到交换机
- channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "");
- // 定义队列的消费者
- DefaultConsumer consumer = new DefaultConsumer(channel) {
- // 获取消息,并且处理,这个方法类似事件监听,如果有消息的时候,会被自动调用
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
- byte[] body) throws IOException {
- // body 即消息体
- String msg = new String(body);
- System.out.println(" [邮件服务] received : " + msg + "!");
- }
- };
- // 监听队列,自动返回完成
- channel.basicConsume(QUEUE_NAME, true, consumer);
- }
- }
我们运行两个消费者,然后发送1条消息:


思考:
1、publish/subscribe与work queues有什么区别。
区别:
1)work queues不用定义交换机,而publish/subscribe需要定义交换机。
2)publish/subscribe的生产方是面向交换机发送消息,work queues的生产方是面向队列发送消息(底层使用默认交换机)。
3)publish/subscribe需要设置队列和交换机的绑定,work queues不需要设置,实际上work queues会将队列绑定到默认的交换机 。
相同点:
所以两者实现的发布/订阅的效果是一样的,多个消费端监听同一个队列不会重复消费消息。
2、实际工作用 publish/subscribe还是work queues。
建议使用 publish/subscribe,发布订阅模式比工作队列模式更强大(也可以做到同一队列竞争),并且发布订阅模式可以指定自己专用的交换机。
④Routing 路由模型(交换机类型:direct)
Routing模型示意图:

P:生产者,向Exchange发送消息,发送消息时,会指定一个routing key。
X:Exchange(交换机),接收生产者的消息,然后把消息递交给 与routing key完全匹配的队列
C1:消费者,其所在队列指定了需要routing key 为 error 的消息
C2:消费者,其所在队列指定了需要routing key 为 info、error、warning 的消息
接下来看代码:
生产者
- public class Send {
- private final static String EXCHANGE_NAME = "test_direct_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明exchange,指定类型为direct
- channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
- // 消息内容,
- String message = "注册成功!请短信回复[T]退订";
- // 发送消息,并且指定routing key 为:sms,只有短信服务能接收到消息
- channel.basicPublish(EXCHANGE_NAME, "sms", null, message.getBytes());
- System.out.println(" [x] Sent '" + message + "'");
- channel.close();
- connection.close();
- }
- }
消费者1
- public class Recv {
- private final static String QUEUE_NAME = "direct_exchange_queue_sms";//短信队列
- private final static String EXCHANGE_NAME = "test_direct_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明队列
- channel.queueDeclare(QUEUE_NAME, false, false, false, null);
- // 绑定队列到交换机,同时指定需要订阅的routing key。可以指定多个
- channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "sms");//指定接收发送方指定routing key为sms的消息
- //channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "email");
- // 定义队列的消费者
- DefaultConsumer consumer = new DefaultConsumer(channel) {
- // 获取消息,并且处理,这个方法类似事件监听,如果有消息的时候,会被自动调用
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
- byte[] body) throws IOException {
- // body 即消息体
- String msg = new String(body);
- System.out.println(" [短信服务] received : " + msg + "!");
- }
- };
- // 监听队列,自动ACK
- channel.basicConsume(QUEUE_NAME, true, consumer);
- }
- }
消费者2
- public class Recv2 {
- private final static String QUEUE_NAME = "direct_exchange_queue_email";//邮件队列
- private final static String EXCHANGE_NAME = "test_direct_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明队列
- channel.queueDeclare(QUEUE_NAME, false, false, false, null);
- // 绑定队列到交换机,同时指定需要订阅的routing key。可以指定多个
- channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "email");//指定接收发送方指定routing key为email的消息
- // 定义队列的消费者
- DefaultConsumer consumer = new DefaultConsumer(channel) {
- // 获取消息,并且处理,这个方法类似事件监听,如果有消息的时候,会被自动调用
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
- byte[] body) throws IOException {
- // body 即消息体
- String msg = new String(body);
- System.out.println(" [邮件服务] received : " + msg + "!");
- }
- };
- // 监听队列,自动ACK
- channel.basicConsume(QUEUE_NAME, true, consumer);
- }
- }
我们发送sms的RoutingKey,发现结果:只有指定短信的消费者1收到消息了

⑤Topics 通配符模式(交换机类型:topics)
Topics模型示意图:

每个消费者监听自己的队列,并且设置带统配符的routingkey,生产者将消息发给broker,由交换机根据routingkey来转发消息到指定的队列。
Routingkey一般都是有一个或者多个单词组成,多个单词之间以“.”分割,例如:inform.sms
通配符规则:
#:匹配一个或多个词
*:匹配不多不少恰好1个词
举例:
audit.#:能够匹配audit.irs.corporate 或者 audit.irs
audit.*:只能匹配audit.irs
从示意图可知,我们将发送所有描述动物的消息。消息将使用由三个字(两个点)组成的Routing key发送。路由关键字中的第一个单词将描述速度,第二个颜色和第三个种类:“
我们创建了三个绑定:Q1绑定了“*.orange.*”,Q2绑定了“.*.*.rabbit”和“lazy.#”。
Q1匹配所有的橙色动物。
Q2匹配关于兔子以及懒惰动物的消息。
下面做个小练习,假如生产者发送如下消息,会进入哪个队列:
quick.orange.rabbit Q1 Q2 routingKey="quick.orange.rabbit"的消息会同时路由到Q1与Q2
lazy.orange.elephant Q1 Q2
quick.orange.fox Q1
lazy.pink.rabbit Q2 (值得注意的是,虽然这个routingKey与Q2的两个bindingKey都匹配,但是只会投递Q2一次)
quick.brown.fox 不匹配任意队列,被丢弃
quick.orange.male.rabbit 不匹配任意队列,被丢弃
orange 不匹配任意队列,被丢弃
下面我们以指定Routing key="quick.orange.rabbit"为例,验证上面的答案
生产者
- public class Send {
- private final static String EXCHANGE_NAME = "test_topic_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明exchange,指定类型为topic
- channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);
- // 消息内容
- String message = "这是一只行动迅速的橙色的兔子";
- // 发送消息,并且指定routing key为:quick.orange.rabbit
- channel.basicPublish(EXCHANGE_NAME, "quick.orange.rabbit", null, message.getBytes());
- System.out.println(" [动物描述:] Sent '" + message + "'");
- channel.close();
- connection.close();
- }
- }
消费者1
- public class Recv {
- private final static String QUEUE_NAME = "topic_exchange_queue_Q1";
- private final static String EXCHANGE_NAME = "test_topic_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明队列
- channel.queueDeclare(QUEUE_NAME, false, false, false, null);
- // 绑定队列到交换机,同时指定需要订阅的routing key。订阅所有的橙色动物
- channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "*.orange.*");
- // 定义队列的消费者
- DefaultConsumer consumer = new DefaultConsumer(channel) {
- // 获取消息,并且处理,这个方法类似事件监听,如果有消息的时候,会被自动调用
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
- byte[] body) throws IOException {
- // body 即消息体
- String msg = new String(body);
- System.out.println(" [消费者1] received : " + msg + "!");
- }
- };
- // 监听队列,自动ACK
- channel.basicConsume(QUEUE_NAME, true, consumer);
- }
- }
消费者2
- public class Recv2 {
- private final static String QUEUE_NAME = "topic_exchange_queue_Q2";
- private final static String EXCHANGE_NAME = "test_topic_exchange";
- public static void main(String[] argv) throws Exception {
- // 获取到连接
- Connection connection = ConnectionUtil.getConnection();
- // 获取通道
- Channel channel = connection.createChannel();
- // 声明队列
- channel.queueDeclare(QUEUE_NAME, false, false, false, null);
- // 绑定队列到交换机,同时指定需要订阅的routing key。订阅关于兔子以及懒惰动物的消息
- channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "*.*.rabbit");
- channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "lazy.#");
- // 定义队列的消费者
- DefaultConsumer consumer = new DefaultConsumer(channel) {
- // 获取消息,并且处理,这个方法类似事件监听,如果有消息的时候,会被自动调用
- public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,
- byte[] body) throws IOException {
- // body 即消息体
- String msg = new String(body);
- System.out.println(" [消费者2] received : " + msg + "!");
- }
- };
- // 监听队列,自动ACK
- channel.basicConsume(QUEUE_NAME, true, consumer);
- }
- }
结果C1、C2是都接收到消息了:


⑥RPC
RPC模型示意图:

基本概念:
Callback queue 回调队列,客户端向服务器发送请求,服务器端处理请求后,将其处理结果保存在一个存储体中。而客户端为了获得处理结果,那么客户在向服务器发送请求时,同时发送一个回调队列地址reply_to。
Correlation id 关联标识,客户端可能会发送多个请求给服务器,当服务器处理完后,客户端无法辨别在回调队列中的响应具体和那个请求时对应的。为了处理这种情况,客户端在发送每个请求时,同时会附带一个独有correlation_id属性,这样客户端在回调队列中根据correlation_id字段的值就可以分辨此响应属于哪个请求。
流程说明:
- 当客户端启动的时候,它创建一个匿名独享的回调队列。
- 在 RPC 请求中,客户端发送带有两个属性的消息:一个是设置回调队列的 reply_to 属性,另一个是设置唯一值的 correlation_id 属性。
- 将请求发送到一个 rpc_queue 队列中。
- 服务器等待请求发送到这个队列中来。当请求出现的时候,它执行他的工作并且将带有执行结果的消息发送给 reply_to 字段指定的队列。
- 客户端等待回调队列里的数据。当有消息出现的时候,它会检查 correlation_id 属性。如果此属性的值与请求匹配,将它返回给应用
分享两道面试题:
面试题:
避免消息堆积?
1) 采用workqueue,多个消费者监听同一队列。
2)接收到消息以后,而是通过线程池,异步消费。
如何避免消息丢失?
1) 消费者的ACK机制。可以防止消费者丢失消息。
但是,如果在消费者消费之前,MQ就宕机了,消息就没了?
2)可以将消息进行持久化。要将消息持久化,前提是:队列、Exchange都持久化
交换机持久化

队列持久化

消息持久化

Spring整合RibbitMQ
下面还是模拟注册服务当用户注册成功后,向短信和邮件服务推送消息的场景
搭建SpringBoot环境
创建两个工程 mq-rabbitmq-producer和mq-rabbitmq-consumer,分别配置1、2、3(第三步本例消费者用注解形式,可以不用配)
1、添加AMQP的启动器:
- <dependency>
- <groupId>org.springframework.bootgroupId>
- <artifactId>spring-boot-starter-amqpartifactId>
- dependency>
- <dependency>
- <groupId>org.springframework.bootgroupId>
- <artifactId>spring‐boot‐starter‐testartifactId>
- dependency>
2、在application.yml中添加RabbitMQ的配置:
- server:
- port: 10086
- spring:
- application:
- name: mq-rabbitmq-producer
- rabbitmq:
- host: 192.168.1.103
- port: 5672
- username: kavito
- password: 123456
- virtualHost: /kavito
- template:
- retry:
- enabled: true
- initial-interval: 10000ms
- max-interval: 300000ms
- multiplier: 2
- exchange: topic.exchange
- publisher-confirms: true
属性说明:
-
template:有关
AmqpTemplate的配置-
retry:失败重试
-
enabled:开启失败重试
-
initial-interval:第一次重试的间隔时长
-
max-interval:最长重试间隔,超过这个间隔将不再重试
-
multiplier:下次重试间隔的倍数,此处是2即下次重试间隔是上次的2倍
-
-
exchange:缺省的交换机名称,此处配置后,发送消息如果不指定交换机就会使用这个
-
-
publisher-confirms:生产者确认机制,确保消息会正确发送,如果发送失败会有错误回执,从而触发重试
当然如果consumer只是接收消息而不发送,就不用配置template相关内容。
3、定义RabbitConfig配置类,配置Exchange、Queue、及绑定交换机。
- public class RabbitmqConfig {
- public static final String QUEUE_EMAIL = "queue_email";//email队列
- public static final String QUEUE_SMS = "queue_sms";//sms队列
- public static final String EXCHANGE_NAME="topic.exchange";//topics类型交换机
- public static final String ROUTINGKEY_EMAIL="topic.#.email.#";
- public static final String ROUTINGKEY_SMS="topic.#.sms.#";
- //声明交换机
- public Exchange exchange(){
- //durable(true) 持久化,mq重启之后交换机还在
- return ExchangeBuilder.topicExchange(EXCHANGE_NAME).durable(true).build();
- }
- //声明email队列
- /*
- * new Queue(QUEUE_EMAIL,true,false,false)
- * durable="true" 持久化 rabbitmq重启的时候不需要创建新的队列
- * auto-delete 表示消息队列没有在使用时将被自动删除 默认是false
- * exclusive 表示该消息队列是否只在当前connection生效,默认是false
- */
- public Queue emailQueue(){
- return new Queue(QUEUE_EMAIL);
- }
- //声明sms队列
- public Queue smsQueue(){
- return new Queue(QUEUE_SMS);
- }
- //ROUTINGKEY_EMAIL队列绑定交换机,指定routingKey
- public Binding bindingEmail(
- return BindingBuilder.bind(queue).to(exchange).with(ROUTINGKEY_EMAIL).noargs();
- }
- //ROUTINGKEY_SMS队列绑定交换机,指定routingKey
- public Binding bindingSMS(
- return BindingBuilder.bind(queue).to(exchange).with(ROUTINGKEY_SMS).noargs();
- }
- }
生产者(mq-rabbitmq-producer)
为了方便测试,我直接把生产者代码放工程测试类:发送routing key是"topic.sms.email"的消息,那么mq-rabbitmq-consumer下那些监听的(与交换机(topic.exchange)绑定,并且订阅的routingkey中匹配了"topic.sms.email"规则的) 队列就会收到消息。
- public class Send {
- RabbitTemplate rabbitTemplate;
- public void sendMsgByTopics(){
- /**
- * 参数:
- * 1、交换机名称
- * 2、routingKey
- * 3、消息内容
- */
- for (int i=0;i<5;i++){
- String message = "恭喜您,注册成功!userid="+i;
- rabbitTemplate.convertAndSend(RabbitmqConfig.EXCHANGE_NAME,"topic.sms.email",message);
- System.out.println(" [x] Sent '" + message + "'");
- }
- }
- }
运行测试类发送5条消息:

web管理界面: 可以看到已经创建了交换机以及queue_email、queue_sms 2个队列,并且向这两个队列分别发送了5条消息


消费者(mq-rabbitmq-consumer)
编写一个监听器组件,通过注解配置消费者队列,以及队列与交换机之间绑定关系。(也可以像生产者那样通过配置类配置)
在SpringAmqp中,对消息的消费者进行了封装和抽象。一个JavaBean的方法,只要添加@RabbitListener注解,就可以成为了一个消费者。
- public class ReceiveHandler {
- //监听邮件队列
- public void rece_email(String msg){
- System.out.println(" [邮件服务] received : " + msg + "!");
- }
- //监听短信队列
- public void rece_sms(String msg){
- System.out.println(" [短信服务] received : " + msg + "!");
- }
- }
属性说明:
-
@Componet:类上的注解,注册到Spring容器 -
@RabbitListener:方法上的注解,声明这个方法是一个消费者方法,需要指定下面的属性:-
bindings:指定绑定关系,可以有多个。值是@QueueBinding的数组。@QueueBinding包含下面属性:-
value:这个消费者关联的队列。值是@Queue,代表一个队列 -
exchange:队列所绑定的交换机,值是@Exchange类型 -
key:队列和交换机绑定的RoutingKey,可指定多个
-
-
启动mq-rabbitmq-comsumer项目

ok,邮件服务和短息服务接收到消息后,就可以各自开展自己的业务了。