学习消息队列Rabbit MQ
#消息队列#记录学习消息队列的知识点。

为什么使用消息队列
有一个支付场景,大家都使用过微信、支付宝支付,比如自动售卖机购买饮料,需要扫码、支付、查询支付结果,出商品。
我们可以发现,售卖机请求支付二维码后,这个请求就结束了,查询支付结果又是另外一个请求,分两种方式,一种是使用http的方式,另外一种是使用推送(消息队列)的方式。不管使用什么方式,我们可以发现,支付的整套流程并不是一次请求就可以搞定。如果一次搞定,高并发的时候服务器的压力就很大,一直在等待支付结果,但用户什么时候支付也不知道。再比如鱼皮哥的做的智能BI项目,AI返回结果的时间会根据用户输入的信息来决定。
消息队列就可以很好解决了同步的问题,采用异步的实现方式。那么消息队列是如何解决同步的问题?
什么是消息队列
存储消息的队列。
- 消息:比如字符串,对象,二进制数据,json等等。
- 队列:先进先出的数据结构。
消息队列的应用场景:
- 耗时的场景(比如支付结果查询)
- 异步化的场景(比如远程控制)
- 应用解耦的场景
- ...
消息队列的好处
- 异步处理
- 削峰填谷
- 应用解耦
消息队列的缺点:
- 要给系统引入额外的中间件,维护成本,资源成本,学习成本。
- 消息队列:就需要考虑,消息丢失,消费消息,数据的一致性。
中间件
- 消息队列也是中间件的一种,我们可以注意到,它是作为生产和消费的中间人,进行传递消息,实现解耦
- 中间件可以简单理解为连接不同系统、应用、网络和数据的一种软件层。中间件将不同系统之间的数据传输与转换进行抽象化和封装,让应用程序只需关注数据流的定义以及操作,而不必关心系统之间的连接细节。
Rabbit MQ是什么
Rabbit MQ是消息队列的一种,生态好,好学习,易于理解,时效性强,支持很多不同语言的客户端,扩展性、可用性都很不错。学习性价比非常高的消息队列,适用于绝大多数中小规模分布式系统。
MQ是消息通信的模型,并发具体实现。现在实现MQ的有两种主流方式:AMQP、JMS。
- JMS限定了必须使用Java语言;AMQP只是协议,不规定实现方式,因此是跨语言的。
- JMS规定了两种消息模型;而AMQP的消息模型更加丰富 ,RabbitMQ是基于AMQP协议,erlang语言开发。
那么接下来我们还需要了解什么才能快速入门?
基本概念
生产者(Publisher)、交换机(Exchange)、路由(Routes)、队列(Queue)、消费者(Consumer)

基本流程
- 生产者会先和rabbit mq建立tcp连接。填写host、username等
- 生产者发送消息给rabbit mq,有Exchange将消息进行路由转发。(如果没有队列需要创建;如果没有交换机需要创建,或者使用默认)
- Exchang将消息路由转发到指定的Queue。
- 消费者会先和rabbit mq建立tcp连接。填写host、username等
- 消费者监听指定的Queue,
- 如果有消息道道Queue,rabbit mq 则将消息推送给消费者
- 消费者接收到消息后,回复确认收到ack。
每一个环节它都是如何实现,怎么进行?这个只是基本的流程,还有更加细节的。但这里我们可以快速入门体验一波。
快速入门
安装
- 官方网站:https://www.rabbitmq.com/
- get started


3. 点击windows install

4.安装rabbit mq之前,还需要安装erlang

5. 安装完erlang后,就可以安装rabbit mq

6. 安装完成后,我们需要启用Rabbit MQ的管理插件(可视化管理界面)
rabbitmq-plugins.bat enable rabbitmq_management

7. services.msc重启mq后生效

8. 访问:http://localhost:15672/,账户名密码都是guest,程序连接端口是5672

Hello World
案例:我们怎么实现,生产一条信息,发送给消费者?

- 引入依赖
- 构建生产者,往队列里面发送消息
(1)首先,需要和rabbit mq建立连接
(2)然后使用Channel,创建出可以进行操作队列、发送消息的工具。
(3)声明队列
(4)发送消息
3. 构建消费者,监听队列面的消息
(1)首先,需要和rabbit mq建立连接
(2)然后使用Channel,创建出可以进行操作队列、发送消息的工具。
(3)声明队列
(4)监听消息
4. 启动生产者
5. 启动消费者

好啦,hello world就结束了。注意:生产者和消费者,创建的队列,必须一摸一样,设置也是。
详细介绍一些状态:

- Name:就是队列名
- Type:表示队列的类型, classic 是先进先出的类型, 当消费者连接到 classic 类型的队列并请求消费消息时,RabbitMQ 将会把消息依次发送给消费者,可以有多个消费者同时从同一个 classic 队列中取消息。classic 队列使用内存存储消息,因此适用于低延迟的场景,但消息的容错能力较差,如果 RabbitMQ 宕机或者重启,未被消费的消息会丢失。
- Features:表示队列的功能特性
- D: 持久化(Durability),指消息队列、交换机和绑定是否持久化存储在磁盘上。如果队列被标记为持久化,那么即使 RabbitMQ 服务器重启,该队列也仍然存在,并且其中的消息也不会丢失。
- TTL:生存时间(Time To Live),指消息在队列中存留的时间。如果一条消息的 TTL 到期后仍未被消费,那么 RabbitMQ 将会自动删除该消息。
- DLX:死信交换机(Dead Letter Exchange),指一种特殊的交换机,它用于处理未被消费的消息。如果一条消息无法被消费,那么 RabbitMQ 将会将其发送到 DLX 中指定的队列中。
- DLK:死信路由键(Dead Letter Routing Key),指用于路由死信消息的路由键。当消息被发送到 DLX 时,RabbitMQ 根据该路由键将消息路由到指定的队列中。
- State: 当前的队列状态以及它是否处于正在使用的状态
- idle:表示队列空闲,没有任何消费者消费。
- running:表示队列正在被使用和消费。
- blocked:表示队列被阻塞,可能出现内存空间不足、磁盘空间不足等问题。
- deleting:表示队列正在被删除,如果队列中仍有未处理的消息,则这些消息将被返回给生产者或转移到 Dead Letter Exchange 中。
- terminated:表示队列已被删除。
- Message:消息的状态
- Ready:指队列中已经准备好可以被消费者消费的消息条数。
- Unacked:指队列中已经被消费者取走但还没有被确认的消息条数。这些未被确认的消息一般是因为消费者出现故障导致,或者消费者在处理消息时还没有进行确认。
- Total:指队列中所有的消息条数,包括已经被消费者取走,但还没有被确认的消息和还没有被消费者取走的消息。
- Message Rate(消息吞吐量)是 RabbitMQ 中的一个性能指标,表示一段时间内消息的接收和处理速率
- Incoming Rate:指消息发送者向 RabbitMQ 发送消息的速率。
- Deliver / Get Rate:指 RabbitMQ 从队列中获取消息并将其推送给消费者或者推送到 Exchange 中的速率。
- Ack Rate:指消费者手动发送 ack 消息来确认消息已被消费的速率。
队列模式
分为work消息模型和publish/subscribe模型。
work消息模型
单向发送(work消息模型)
- 指的是1对1
- 生产者队列发送消息,消费者接到队列路由来的消息。

就是用hello world的例子。这里不过多介绍
多消费者(work消息模型)
- 指的是1对多
- 生产者队列发送消息,多个消费者竞争消息。
谁抢到,谁执行,work消息模型,竞争消费者模式。

其实和单向发送的写法一样,只不过多了一个消费者
publish/subscribe模型
Work 消息模型:在这种模型中,消息被发送到一个队列中,多个消费者从该队列中接收消息并按照一定的顺序进行处理。消息只能被一个消费者接收 。如果我们想要一条消息被多个消费者消费怎么办呢?就需要使用到我们的 Publish/Subscribe 模型 ,订阅模型。

可以看到,之前生产者直接对接队列

- 现在变成了,生成者对接交换机,交换机下发到对应的队列。通过路由把消息转发到不同的队列上,把交换机和队列关联起来。
- 绑定规则是什么,规则就是有多种模式,交换机有多少类别:fanout\direct\topic
交换机做了什么呢?
- 接收生产者发送的消息。另一方面:知道如何处理消息,例如递交给某个特别队列、递交给所有队列、或是将消息丢弃。
- 只负责转发消息,不具备存储消息的能力,因此如果没有任何队列与Exchange绑定,或者没有符合路由规则的队列,那么给交换机发送消息,消息会丢失!
- 使用了交换机,生产者不在声明Queue,发送消息给Exchange,不在发送到Queue
fanout:广播模型

特点:消息会被路由到所有绑定到该交换机的队列上。
举个场景:有10台自动售卖机,我们想给他推送广告,那么我们就可以使用这样的交换机来推送。
接下来我们看一下如何实现的。
- 创建了两个队列。绑定同一个交换机

Direct:定向模型
绑定:可以让交换机,发送消息给某个队列,上面是全部发,这里是指定发,通过routingKey路由键。也就是,交换机也可以单独发送到指定的队列。
绑定关系:完全匹配字符串

P:生产者,向Exchange发送消息,发送消息时,会指定一个routing key。
X:Exchange(交换机),接收生产者的消息,然后把消息递交给 与routing key完全匹配的队列
C1:消费者,其所在队列指定了需要routing key 为 orange 的消息
C2:消费者,其所在队列指定了需要routing key 为 black、green的消息
可以看到,我们这里使用上了路由键。xiaowang,xiaoli


疑问:如果路由键不指定呢??
- 如果生产者在发送消息时未指定 routing key,那么消息到达 Exchange 后无法判断应该转发到哪个队列,而且也没有任何与 Exchange 绑定的队列可以接收这个消息。因此,在 Direct exchange 模型中,如果生产者未指定 routing key,那么消息将无法被消费者收到。
Topic:通配符
消息会根据一个 模糊的 路由键转发到指定的队列
绑定关系:可以模糊匹配多个绑定·
- *:匹配一个单词,比如*.orange,那么 a.orange, b.orange 都能匹配·
- #:匹配 0个或多个单词,比如 a.#,那么 a.a, a.b, a.a.a 都能匹配


核心特性
消息过期机制
可以给每条消息设置一个有效期,一段时间内未被消费者处理就过期了
场景:比如消费者没有同网络,那么这个时候随意发送的消息就会积累过多,其实这些消息当时如果没有用,后面就不需要了,那么可以设置一个有效期。分两种:给这个队列所有的消息指定过期时间和给指定消息过期时间:
队列所有消息指定过期时间
在队列声明的时候,需要声明过期时间参数的设置。
如果在过期时间内,还没有消费者取消息,消息才会过期。
注意,如果消息已经接收到,但是没确认,是不会过期的。就会出现如下这种情况,收到了但是没有确认。

如果消息处于待消费状态并且过期时间到达后,消息将被标记为过期。但是,如果消息已经被消费者消费,并且正在处理过程中,即使过期时间到达,消息仍然会被正常处理。如果重启消费者,消息也不会重新消费,直接没了???真的吗?
给某条消息指定过期时间
在发消息的时候,指定过期时间
消息确认机制
为什么需要消息确认机制呢? 它确保了消息的可靠性传递和处理 。如果你发送了一条消息,比如消费者在处理过程中,出现异常,那么如果你还没确认消息接收到,消费者重新起来的时候,就可以又继续消费。服务器也知道消息被消费或其他情况,也就是需要给我一个反馈。反馈方式有三种:
- ack:消费成功,如果配置autoack,那么消费者一收到消息,会自动执行ack。
- nack:消费失败
- reject:拒绝
代码实现:
第一个参数:标识消息的唯一编号,用于在手动确认消息时指定要确认的消息。
第二个参数:批量确认,表示是否一次性确认所有历史消息,包括当前这一条。false表示不批量确认
第三个参数:表示是否重新入队,可用于重试。false表示不进行
- `basicNack()` 方法可以批量拒绝多个消息并选择是否重新入队;
- `basicReject()` 方法只能拒绝单个消息并选择是否重新入队。
死信队列
死信:过期的消息、拒收的消息、消息队列满了、处理失败的消息的统称
死信队列,死信交换机:都是普通的交换机,和我们平常使用的没什么区别,只是起了个名字
Spring整合RibbitMQ
- 引入依赖
- 添加rabbit mq的配置
- 创建交换机,创建队列,绑定
- 创建生产者,发送消息
- 创建消费者,监听消费消息。
好了,分享就到这里。
