系列说明
本系列主要讲解RabbitMQ,讲解其特性,例如消息持久化、消息TTL、消息的优先、延迟消息、消息可靠性、消费模式以及在Spring Boot中使用RabbitMQ,代码在我的Github上
RabbitMQ介绍
RabbitMQ使用Erlang语言开发基于AQMP协议的开源消息队列,RabbitMQ主要有以下特点:
- 高可靠性: RabbitMQ依靠消息确认、持久化等实现高可靠性,但其吞吐量不太高
- 高可用: RabbitMQ支持分布式部署,多个RabbitMQ服务器组成一个集群形成一个逻辑Broker
- 多种协议支持: RabbitMQ基于AQMP协议,但是可以通过安装插件支持其它协议,例如STOMP、MQTT协议等
- 多种客户端语言支持: RabbitMQ提供Java、C++等多种客户端语言支持
- 管理页面: RabbitMQ提供Web管理页面以便可视化管理
AQMP
RabbitMQ基于AQMP协议开发的消息队列,AQMP协议在之前消息队列(一)中已经简单的介绍了,这里就简单的介绍一下:

- Broker: Broker指的是代理消息队列,是一个逻辑概念,指的是RabbitMQ服务器,其可以有多个Vritual Host组成。
- Virtual Host: Vritual Host是一个虚拟概念,类似于权限控制组,一个Vritual Host可以有若干Exchange和Queue,权限控制的最小粒度是Virtual Host
- Exchange: Exchange叫交换机,其可以多个Queue根据路由规则(Routing Key)绑定。Exchange接收生产者发送的消息,根据其类型(ExchangeType)和路由规则(Routing Key)把消息发送给队列。
- Bingding: Binding联系Exchange和Queue
- Connection: Connection在RabbitMQ中是一个客户端和Broker之间的TCP连接
- Channel: Channel在RabbitMQ中叫做信道,有Connection创建,并且一个Connection可以创建多个Channel。在RabbitMQ必须通过Channel才能发送消息。之所以需要Channel,主要因为TCP连接过于昂贵。
需要注意的地方:
- 如果消息发送到Exchange后,Exchange不能通过路由规则找到合适的队列,该消息将会被删除
- RabbitMQ建议客户端线程之间不要共用Channel,而是共用Connection
RabbitMQ使用
RabbitMQ-Java官方提供了简单的使用教程,这里就简单的提一下,具体可见其网友翻译版本:RabbitMQ入门教程
这里展示的是RabbitMQ发送消息
1 | public class Sender { |
这里展示使用RabbitMQ接收消息
1 | public class Receiver { |
RabbitMQ的消息传递模型
RabbitMQ主要通过ExchangeType来设置消息传递模型,主要有下面4种模型,其中Header模型用的少:
- Direct模型
- Fanout模型
- Topic模型
- Header模型
Direct模型
Direct模型顾名思义指的是直接连接,只有当消息中的Routing Key与Queue绑定到Exchange的Routing Key一致,才会转发消息给该Queue
Fanout模型
Fanout模型类似于订阅/发布模型,Exchange会把消息转发给所有绑定到该Exchange上的Queue
Topic模型
Topic模型类型与Servlet的URL匹配模型,其会匹配消息的Routing Key和Queue绑定到Exchange的Routing Key,使用通配符匹配。有#和*两种通配符,#代表0个或多个字符,*代表1个字符
RabbitMQ的持久化
首先RabbitMQ的持久化是异步持久化模型,也就是说在特定情况下,可能造成消息丢失。比如在RabbitMQ Server回调RabbitMQ Producer Client的接口表明已经接收到该消息,但是由于是异步持久化可能还没有把消息持久化到磁盘中,这时候MQ-Server断电就会导致消息的丢失
RabbitMQ中消息的持久化需要保证Exchange、Queue、Message都进行持久化操作。需要注意的是:Exchange、Queue的声明时幂等的。幂等指说多次声明产生的结果都是一样,也就是说如果其不存在则创建,存在则返回且不会对其产生任何影响,但是如果声明已存在的队列,且其属性不同则会抛出异常。
Exchange的持久化
RabbitMQ声明Exchange有几种方法,但主要使用下面方法,其中第三个参数表示是否将该Exchange持久化
1 | /** |
Queue的持久化
RabbitMQ声明Queue与Exchange的方法类型,同样使用durable参数表示是否将该Queue进行持久化操作,下面是其中一个方法
1 | /** |
部分参数说明:
- exclusive:表明该队列是否是排它队列。如果一个队列被声明为排他队列,该队列仅对首次申明它的连接可见,并在连接断开时自动删除。这里需要注意三点:1. 排他队列是基于连接可见的,同一连接的不同信道是可以同时访问同一连接创建的排他队列;2.“首次”,如果一个连接已经声明了一个排他队列,其他连接是不允许建立同名的排他队列的,这个与普通队列不同;3.即使该队列是持久化的,一旦连接关闭或者客户端退出,该排他队列都会被自动删除的,这种队列适用于一个客户端发送读取消息的应用场景。
- autoDelete:表明该队列是否自动删除。自动删除,如果该队列没有任何订阅的消费者的话,该队列会被自动删除。这种队列适用于临时队列。
Message的持久化
消息的持久化需要在生产者发送消息时设置消息属性,以表明该消息时持久化消息。下面是消息发送的一个API
1 | /** |
部分参数说明:
- exchange:表示该消息发送到哪个Exchange
- routingKey:表示该消息的Routing Key
- props:表示该消息的属性
- body:消息实体
BasicProperties定义如下:
1 | public BasicProperties( |
其中RabbitMQ提供了属性实现已经更简单的配置消息属性:
1 | /** Empty basic properties, with no fields set */ |
当然可以使用时编程自定义设置消息属性:
1 | AMQP.BasicProperties.Builder builder = new Builder(); |
消息什么时候刷到磁盘(持久化)
写入文件前会有一个Buffer,大小为1M,数据在写入文件时,首先会写入到这个Buffer,如果Buffer已满,则会将Buffer写入到文件(未必刷到磁盘)。
有个固定的刷盘时间:25ms,也就是不管Buffer满不满,每个25ms,Buffer里的数据及未刷新到磁盘的文件内容必定会刷到磁盘。
每次消息写入后,如果没有后续写入请求,则会直接将已写入的消息刷到磁盘:使用Erlang的receive x after 0实现,只要进程的信箱里没有消息,则产生一个timeout消息,而timeout会触发刷盘操作。
TTL
TTL(Time To Live)表示存活时间。RabbitMQ中可以对Queue和Message设置TTL,以控制Queue和Message的存活时间。
Queue TTL
队列的存活时间指的是Queue在自动删除前可以处于未使用状态的时间。未使用状态指的是Queue上没有Consumer、Queue没有被重新声明。队列的存活时间在队列第一次声明时通过指定队列的属性”x-expires”指定,单位是毫秒,代码如下:
1 | Map<String, Object> queueArgs = new HashMap<>(); |
Message TTL
消息的存活时间指的是消息在队列中的存活时间,超过该时间消息将被删除或者不能传递给消费者。消息的存活时间可以通过设置每条消息的存活时间或者设置某条队列中的所以存活时间,当两者都有时,时间小的有效。
设置消息属性
针对每条消息可以在发送消息时设置消息属性
1 | // 设置消息属性-TTL为30s |
设置队列属性
通过设置队列中消息的TTL属性,然后传入该队列的所有消息都有该TTL属性
1 | Map<String, Object> queueArgs = new HashMap<>(); |
Reference
https://www.jianshu.com/p/64357bf35808?utm_campaign=maleskine&utm_content=note&utm_medium=seo_notes&utm_source=recommendation
http://blog.csdn.net/wanbf123/article/details/78052419
http://blog.csdn.net/u013256816/article/details/60875666
http://blog.csdn.net/u013256816/article/details/54916011