一、走进RabbitMQ

1.1 消息中间件介绍

消息中间件(消息队列)是分布式系统中重要的组件,主要解决应用耦合,异步消 息,流量削锋等问题实现高性能,高可用,可伸缩和最终一致性[架构] 使用较多的消息 队列有ActiveMQ,RabbitMQ,ZeroMQ,Kafka,MetaMQ,RocketMQ
以下介绍消息队列在实际应用中常用的使用场景:异步处理,应用解耦,流量削锋和消
息通讯四个场景

1.2 什么是RabbitMQ

RabbitMQ 是一个由 Erlang 语言开发的 AMQP 的开源实现。
AMQP :Advanced Message Queue,高级消息队列协议。它是应用层协议的一个开放 标准,为面向消息的中间件设计,基于此协议的客户端与消息中间件可传递消息,并不 受产品、开发语言等条件的限制。
RabbitMQ 最初起源于金融系统,用于在分布式系统中存储转发消息,在易用性、扩展 性、高可用性等方面表现不俗。具体特点包括:

1.可靠性(Reliability)
RabbitMQ 使用一些机制来保证可靠性,如持久化、传输确认、发布确认。

2.灵活的路由(Flexible Routing)
在消息进入队列之前,通过 Exchange 来路由消息的。对于典型的路由功能,RabbitMQ 已经提供了一些内置的 Exchange 来实现。针对更复杂的路由功能,可以将多个 Exchange 绑定在一起,也通过插件机制实现自己的 Exchange

3.消息集群(Clustering)
多个 RabbitMQ 服务器可以组成一个集群,形成一个逻辑 Broker 。

4.高可用(Highly Available Queues)
队列可以在集群中的机器上进行镜像,使得在部分节点出问题的情况下队列仍然可用。

5.多种协议(Multi-protocol)
RabbitMQ 支持多种消息队列协议,比如 STOMP、MQTT 等等。

6.多语言客户端(Many Clients)
RabbitMQ 几乎支持所有常用语言,比如 Java、.NET、Ruby 等等。

7.管理界面(Management UI)
RabbitMQ 提供了一个易用的用户界面,使得用户可以监控和管理消息 Broker 的许多方 面。

8.跟踪机制(Tracing)
如果消息异常,RabbitMQ 提供了消息跟踪机制,使用者可以找出发生了什么。

9.插件机制(Plugin System)
RabbitMQ 提供了许多插件,来从多方面进行扩展,也可以编写自己的插件。

1.3 主要概念

RabbitMQ Server: 也叫broker server,它是一种传输服务。 他的角色就是维护一条从Producer到Consumer的路线,保证数据能够按照指定的方式进行传输。

Producer (Publisher): 消息生产者,如图A、B、C,数据的发送方。消息生产者连接RabbitMQ服
务器然后将消息投递到Exchange。

Consumer:消息消费者,如图1、2、3,数据的接收方。消息消费者订阅队列,
RabbitMQ将Queue中的消息发送到消息消费者。

Exchange:生产者将消息发送到Exchange(交换器),由Exchange将消息路由到一个 或多个Queue中(或者丢弃)。Exchange并不存储消息。RabbitMQ中的Exchange有 direct、fanout、topic、headers四种类型,每种类型对应不同的路由规则。

Queue:(队列)是RabbitMQ的内部对象,用于存储消息。消息消费者就是通过订阅 队列来获取消息的,RabbitMQ中的消息都只能存储在Queue中,生产者生产消息并最终 投递到Queue中,消费者可以从Queue中获取消息并消费。多个消费者可以订阅同一个 Queue,这时Queue中的消息会被平均分摊给多个消费者进行处理,而不是每个消费者 都收到所有的消息并处理。

RoutingKey:生产者在将消息发送给Exchange的时候,一般会指定一个routing key, 来指定这个消息的路由规则,而这个routing key需要与Exchange Type及binding key联 合使用才能最终生效。在Exchange Type与binding key固定的情况下(在正常使用时一 般这些内容都是固定配置好的),我们的生产者就可以在发送消息给Exchange时,通过 指定routing key来决定消息流向哪里。RabbitMQ为routing key设定的长度限制为255 bytes。

Connection: (连接):Producer和Consumer都是通过TCP连接到RabbitMQ Server 的。以后我们可以看到,程序的起始处就是建立这个TCP连接。

Channels: (信道):它建立在上述的TCP连接中。数据流动都是在Channel中进行 的。也就是说,一般情况是程序起始建立TCP连接,第二步就是建立这个Channel。

VirtualHost:权限控制的基本单位,一个VirtualHost里面有若干Exchange和 MessageQueue,以及指定被哪些user使用。
image.png

生产者与Broker,消费者与Broker都只会建立一条连接。在这个连接中会开辟很多条信道(channel),当连接中断时,broker会监控到中断并将消息存储起来,避免了消息丢失。

二、Docker 安装RabbitMQ

  1. docker pull rabbitmq:3.8.2-management
  2. docker run -d --name rabbitmq -p 5671:5671 -p 5672:5672 -p 4369:4369 -p 25672:25672 -p 15671:15671 -p 15672:15672 rabbitmq:3.8.2-management

4369,25672(Erlang发现&集群端口)
5672,5672(AMQP端口)
15672(web管理后台端口)
61613,61614(STOMP协议端口)
1883,8883(MQTT协议端口)

https://www.rabbitmq.com/networking.html

image.png

  1. # 设置自动重启
  2. docker update rabbitmq --restart=always

http://localhost:15672/
image.png
默认账户密码:guest / guest
可以根据环境配置不同的虚拟主机,进行环境的隔离
image.png

三、RabbitMQ运行机制

image.png

3.1 Exchange类型

Exchange分发消息时根据类型的不同分发策略有区别,目前共四种类型:direct、fanout、topic、headers。headers匹配AMQP消息的header而不是路键,headers交换机和direct交换机完全一致,但性能差很多,目前几乎用不到了,所以直接看另外三种类型,其中direct和header都是点对点的实现,而fanout和topic都是发布订阅的实现,header效率低下已经弃用了。

  • Direct(点对点通信模式)

image.png

  • Fanout (广播模式)

image.png

  • Topic (主题发布订阅模式)

image.png

3.2 总结

生产者发送消息给交换机,消费者监听的是队列,交换机把消息交给队列

四、RabbitMQ使用

4.1 创建交换机

image.png

  • durability:是否持久化,
    • 参数:transient零时的,当服务重启时,交换机将会消失
  • auto delete:是否自动删除
    • yes:交换机没有绑定任何东西时,将会自动删除自己
  • Internal:是否是内部交换机

    • yes:将不能转发消息,仅供内部转发路由使用

      4.2 创建队列

      image.png
  • durability:是否持久化,

    • 参数:transient零时的,当服务重启时,交换机将会消失
  • auto delete:是否自动删除

    • yes: 没有连接时就会自动删除自身。

      4.3 交换机绑定队列

      image.png

      4.4 测试

      4.4.1 direct 直接交换机

      1、准备四个队列,一个直接交换机,并进行绑定
      image.png
      2、发送消息
      image.png
      3、队列收到消息
      image.png
      4、点击队列后查看消息
      image.png
  • ack mode:

    • Nack message requeue true: get消息并再次放入队列
    • Ack message requeue true:get消息并删除

      4.4.2 fanout 扇出交换机

      image.pngimage.png
      发送消息后,无论写不写路由键(route key),所有被绑定的队列都会收到消息。

      4.4.3 topic 主题发布订阅交换机

      image.png
      发送消息,smiler.news;根据路由key,可以预计四个队列都会收到消息。
      因为:smiler.news会匹配#.news和smiler.#
      image.png
      四个队列,收到消息
      image.png

      五、Springboot整合RabbitMQ

      1、引入 spring-boot-starter-amqp 2、application.yml配置 3、测试RabbitMQ

      1. 3.1 AmqpAdmin:管理组件
      2. 3.2 RabbitTemplate:消息发送处理组件

5.1 RabbitAutoConfiguration

image.png

5.2 SpringBoot配合RabbitMQ使用

  1. spring.rabbitmq.host=127.0.0.1
  2. spring.rabbitmq.port=5672
  3. spring.rabbitmq.virtual-host=/

5.2.1 AmqpAdmin

  1. /**
  2. * 1、如何创建Exchange、Queue、Binding
  3. * 1)、使用AmqpAdmin进行创建
  4. * 2、如何收发消息
  5. */
  6. @Test
  7. public void createExchange() {
  8. // DirectExchange(String name, booean durable, boolean autoDelete, Map<String, Object> arguments)
  9. DirectExchange directExchange = new DirectExchange(CHANGE_NAME, true, false);
  10. amqpAdmin.declareExchange(directExchange);
  11. LOGGER.info("Exchange[{}]创建成功", CHANGE_NAME);
  12. }
  13. @Test
  14. public void createQueue() {
  15. // (String name, boolean durable, boolean exclusive, boolean autoDelete, Map<String, Object> arguments)
  16. Queue queue = new Queue(QUEUE_NAME, true, false, false);
  17. amqpAdmin.declareQueue(queue);
  18. LOGGER.info("Queue[{}]创建成功", QUEUE_NAME);
  19. }
  20. @Test
  21. public void createBinding() {
  22. // (String destination, DestinationType destinationType, String exchange, String routingKey,Map<String, Object> arguments)
  23. // 将exchange指定的交换机和destination目的地进行绑定,使用routingkey作为指定的路由键
  24. Binding binding = new Binding(QUEUE_NAME, Binding.DestinationType.QUEUE, CHANGE_NAME, ROUTING_KEY, null);
  25. amqpAdmin.declareBinding(binding);
  26. LOGGER.info("Binding[{}]创建成功", BINDING_NAME);
  27. }

5.2.2 RabbitTemplate

  1. @Test
  2. public void sendMessageTest() {
  3. String msg = "hello world";
  4. rabbitTemplate.convertAndSend(CHANGE_NAME, ROUTING_KEY, msg);
  5. LOGGER.info("消息发送完成{}", msg);
  6. }

可以配置MyRabbitConfig控制发送消息的格式

  1. @Configuration
  2. public class MyRabbitConfig {
  3. /**
  4. * 将rabbitMQ发送的消息对象转化为json格式
  5. * @return
  6. */
  7. @Bean
  8. public MessageConverter messageConverter(){
  9. return new Jackson2JsonMessageConverter();
  10. }
  11. }

5.3 RabbitListener

可以标注在类或者方法上

  1. **
  2. * @description: RabbitHandler 根据消息体类型接收消息
  3. * @Author: wangchao
  4. * @Date: 2021/4/24
  5. */
  6. @Service
  7. @RabbitListener(queues = "hello-java-queue")
  8. public class RabbitHandler {
  9. @org.springframework.amqp.rabbit.annotation.RabbitHandler
  10. public void receiveMessage(Message message, String content) {
  11. System.out.println("111消息内容:" + content);
  12. }
  13. @org.springframework.amqp.rabbit.annotation.RabbitHandler
  14. public void receiveMessage(Message message, List<String> content) {
  15. System.out.println("222消息内容:" + content);
  16. }
  17. }

5.4 RabbitHandler

作用在方法上

  1. /**
  2. * @description: RabbitListener 可以标注在类上,也可以标注在方法上
  3. * @Author: wangchao
  4. * @Date: 2021/4/24
  5. */
  6. @Service
  7. public class RabbitMqListener {
  8. /**
  9. * queues:声明需要监听的所有队列
  10. * org.springframework.amqp.core.Message
  11. * 参数可以写一下类型
  12. * 1、Message message:原生消息详细信息。头+体
  13. * 2、T<消息发送类型> String content
  14. * 3、Channel channel:
  15. *
  16. * Queue:可以很多人都来监听。只要收到消息,队列删除消息,而且只能有一个人收到消息
  17. * 场景:
  18. * 1)服务启动多个,只能有一个客户端收到消息
  19. * 2)只有一个消息完全处理完成,才会接收下一个消息
  20. * @param message
  21. * @param content
  22. */
  23. @RabbitListener(queues = "hello-java-queue")
  24. public void receiveMessage1(Message message,String content) {
  25. System.out.println("接收到消息。。。。内容:" + message + "\n===>类型:" + message.getClass());
  26. System.out.println("===>body:"+ Arrays.toString(message.getBody()));
  27. MessageProperties messageProperties = message.getMessageProperties();
  28. System.out.println("===>messageProperties:"+ JSON.toJSONString(messageProperties));
  29. System.out.println("===>content:"+content);
  30. }
  31. @RabbitListener(queues = "hello-java-queue")
  32. public void receiveMessage2(Message message,String content) {
  33. System.out.println("===>222content:"+content);
  34. }
  35. @RabbitListener(queues = "hello-java-queue")
  36. public void receiveMessage3(Message message,String content) {
  37. System.out.println("===>333content:"+content);
  38. }
  39. }

六、RabbitMQ消息确认机制-可靠抵达

  • 保证消息不丢失,可靠抵达,可以使用事务消息,性能下降250倍,为此引入确认机制
  • publisher comfirmCallback 确认模式
  • publish returnCallback 未投递到queue退回模式
  • consumer ack机制

image.png

  1. /**
  2. * 定制rabbitTemplate
  3. * 1、服务器收到消息就会调
  4. * 1.1、spring.rabbitmq.publisher-confirm-type=correlated
  5. * 1.2、设置确认会调
  6. * 2、消息正确抵达队列进行会调
  7. */
  8. @PostConstruct // MyRabbitConfig对象构造器创建完成之后再调用这个方法
  9. public void initRabbitTemplate() {
  10. rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
  11. /**
  12. * 1、只要消息抵达broker,ack就是true
  13. * @param correlationData 当前消息的唯一关联数据(这个是消息的唯一id)
  14. * @param ack 消息是否成功收到
  15. * @param cause 失败的原因
  16. */
  17. @Override
  18. public void confirm(CorrelationData correlationData, boolean ack, String cause) {
  19. System.out.println("confirm...correlationData[" + correlationData + "]");
  20. System.out.println("confirm...ack[" + ack + "]");
  21. System.out.println("confirm...cause[" + cause + "]");
  22. }
  23. });
  24. }

image.png

  1. // 设置消息抵达队列的确认回调
  2. rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
  3. /**
  4. * 只要消息没有投递给指定的队列,就会触发这个失败的回调
  5. * @param message 投递失败的消息详细信息
  6. * @param replyCode 回复的状态码
  7. * @param replyText 回复的文本内容
  8. * @param exchange 当时这个消息发送给哪个交换机
  9. * @param routingKey 当时这个消息用哪个路由键
  10. */
  11. @Override
  12. public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
  13. System.out.println("Fail Message[" + message + "]");
  14. System.out.println("Fail Message[" + replyCode + "]");
  15. System.out.println("Fail Message[" + replyText + "]");
  16. System.out.println("Fail Message[" + exchange + "]");
  17. System.out.println("Fail Message[" + routingKey + "]");
  18. }
  19. });


image.png

  1. @RabbitListener(queues = "hello-java-queue")
  2. public void receiveMessage3(Message message, String content, Channel channel) {
  3. long deliveryTag = message.getMessageProperties().getDeliveryTag();
  4. // 确认,非批量模式
  5. try {
  6. channel.basicAck(deliveryTag,false);
  7. } catch (IOException e) {
  8. e.printStackTrace();
  9. }
  10. System.out.println("===>333content:" + content);
  11. }
  12. @RabbitListener(queues = "hello-java-queue")
  13. public void receiveMessage4(Message message, String content, Channel channel) {
  14. long deliveryTag = message.getMessageProperties().getDeliveryTag();
  15. // 确认,非批量模式
  16. try {
  17. // requeue:丢弃 requeue=true 发回服务器,服务器重新入队。
  18. channel.basicNack(deliveryTag,false,false);
  19. // long deliveryTag, boolean requeue
  20. // channel.basicReject();
  21. } catch (IOException e) {
  22. e.printStackTrace();
  23. }
  24. System.out.println("===>333content:" + content);
  25. }
  1. # 收到消息后进行手动确认(默认是自动确认的)
  2. spring.rabbitmq.listener.simple.acknowledge-mode=manual
  3. 3、消费端确认(保证每个消费被准确消费,此时才可以broker删除这个消息)。
  4. * 3.1、默认是自动确认的,只要消息接收到,服务端就会移除这个消息。这会带来一个问题:我们收到很多消息,自动回复给服务器ack,只有一个消息处理
  5. * 成功,宕机了。发生消息丢失;
  6. * 开启手动确认模式之后,只要我们没有明确告诉MQ签收,没有ack,消息就一直是unack状态。即使consumer宕机。消息也不会丢失,会重新变为ready,
  7. * 下次有新的consumer进来就发给他。
  8. * 3.2、如何签收
  9. * channel.basicAck(deliveryTag,false);签收,业务成功完成就签收
  10. * channel.basicNack(deliveryTag,false,false);拒签,业务失败,拒签