一、RabbitMQ 概述
1 消息队列
消息(Message)是指在应用间传送的数据。消息可以非常简单,比如只包含文本字符串,也可以更复杂,可能包含嵌入对象。不建议传递对象,如果需要传递复杂数据建议传递Json。
消息队列(Message Queue)是一种应用间的通信方式,消息发送后可以立即返回,由消息系统来确保消息的可靠传递。消息发布者只管把消息发布到 MQ 中而不用管谁来取,消息使用者只管从 MQ 中取消息而不管是谁发布的。这样发布者和使用者都不用知道对方的存在。
消息队列用于业务解耦、最终一致性、广播、错峰流控等等情况
2 RabbitMQ 特点
RabbitMQ 是一个由 Erlang 语言开发的 AMQP 的开源实现。
AMQP :Advanced Message Queue,高级消息队列协议。它是应用层协议的一个开放标准,为面向消息的中间件设计,基于此协议的客户端与消息中间件可传递消息,并不受产品、开发语言等条件的限制。
RabbitMQ 最初起源于金融系统,用于在分布式系统中存储转发消息,在易用性、扩展性、高可用性等方面表现不俗。具体特点包括:
- 可靠性(Reliability):RabbitMQ 使用一些机制来保证可靠性,如持久化、传输确认、发布确认。
- 灵活的路由(Flexible Routing):在消息进入队列之前,通过 Exchange 来路由消息的。对于典型的路由功能,RabbitMQ 已经提供了一些内置的 Exchange 来实现。针对更复杂的路由功能,可以将多个 Exchange 绑定在一起,也通过插件机制实现自己的 Exchange。
- 消息集群(Clustering):多个 RabbitMQ 服务器可以组成一个集群,形成一个逻辑 Broker。高可用队列可以在集群中的机器上进行镜像,使得在部分节点出问题的情况下队列仍然可用。
- 多种协议(Multi-protocol):RabbitMQ 支持多种消息队列协议,比如 STOMP、MQTT 等等。
- 多语言客户端(Many Clients):RabbitMQ 几乎支持所有常用语言,比如 Java、.NET、Ruby 等等。
- 管理界面(Management UI):RabbitMQ 提供了一个易用的用户界面,使得用户可以监控和管理消息 Broker 的许多方面。
- 跟踪机制(Tracing):如果消息异常,RabbitMQ 提供了消息跟踪机制,使用者可以找出发生了什么。
二 RabbitMQ的消息发送和接收机制
1 概述
所有 MQ 产品从模型抽象上来说都是一样的过程: 消费者(consumer)订阅某个队列。生产者(producer)创建消息,然后发布到队列(queue)中, 最后将消息发送到监听的消费者。
RabbitMQ的内部接收如下:
- Message:消息,消息是不具体的,它由消息头和消息体组成。消息体是不透明的,而消息头则由一系列的可选 属性组成,这些属性包括routing-key(路由键)、priority(相对于其他消息的优先权)、deliverymode(指出该消息可能需要持久性存储)等。
- Publisher:消息的生产者,也是一个向交换器发布消息的客户端应用程序。
- Exchange:交换器,用来接收生产者发送的消息并将这些消息路由给服务器中的队列。
- Binding:绑定,用于消息队列和交换器之间的关联。一个绑定就是基于路由键将交换器和消息队列连接起来的路 由规则,所以可以将交换器理解成一个由绑定构成的路由表。
- Queue:消息队列,用来保存消息直到发送给消费者。它是消息的容器,也是消息的终点。一个消息可投入一个 或多个队列。消息一直在队列里面,等待消费者连接到这个队列将其取走。
- Connection:网络连接,比如一个TCP连接。
- Channel:信道,多路复用连接中的一条独立的双向数据流通道。信道是建立在真实的TCP连接内地虚拟连接, AMQP 命令都是通过信道发出去的,不管是发布消息、订阅队列还是接收消息,这些动作都是通过信道 完成。因为对于操作系统来说建立和销毁 TCP 都是非常昂贵的开销,所以引入了信道的概念,以复用一 条 TCP 连接。
- Consumer:消息的消费者,表示一个从消息队列中取得消息的客户端应用程序。
- Virtual Host:虚拟主机,表示一批交换器、消息队列和相关对象。虚拟主机是共享相同的身份认证和加密环境的独立 服务器域。每个 vhost 本质上就是一个 mini 版的 RabbitMQ 服务器,拥有自己的队列、交换器、绑定 和权限机制。vhost 是 AMQP 概念的基础,必须在连接时指定,RabbitMQ 默认的 vhost 是 / 。
- Broker:表示消息队列服务器实体。
2 AMQP 中的消息路由
AMQP 中消息的路由过程和 Java 开发者熟悉的 JMS 存在一些差别,AMQP 中增加了 Exchange 和 Binding 的角色。生产者把消息发布到 Exchange 上,消息最终到达队列并被消费者接收,而 Binding 决定交换器的消息应该发送到那个队列
3 Exchange 类型
Exchange分发消息时根据类型的不同分发策略有区别,目前共四种类型:direct、fanout、topic、headers 。
- headers:匹配 AMQP 消息的 header 而不是路由键,此外 headers 交换器和 direct 交换器完全一致,但性能差很多,目前几乎用不到了
- direct:消息中的路由键(routing key)如果和 Binding 中的 binding key 一致, 交换器就将消息发到对应的 队列中。路由键与队列名完全匹配,如果一个队列绑定到交换机要求路由键为“dog”,则只转发 routing key 标记为“dog”的消息,不会转发“dog.puppy”,也不会转发“dog.guard”等等。它是完全匹配、单播的模式。
- fanout:每个发到 fanout 类型交换器的消息都会分到所有绑定的队列上去。fanout 交换器不处理路由键,只是简单的将队列绑定到交换器上,每个发送到交换器的消息都会被转发到与该交换器绑定的所有队列上。 很像子网广播,每台子网内的主机都获得了一份复制的消息。fanout 类型转发消息是最快的。
- topic:topic 交换器通过模式匹配分配消息的路由键属性,将路由键和某个模式进行匹配,此时队列需要绑定 到一个模式上。它将路由键和绑定键的字符串切分成单词,这些单词之间用点隔开。它同样也会识别两个通配符:符号“#”和符号“*”。#匹配0个或多个单词,*匹配不多不少一个单词。
三、Java RabbitMQ
1 基础操作
依赖
<dependencies><dependency><groupId>com.rabbitmq</groupId><artifactId>amqp-client</artifactId><version>5.1.1</version></dependency></dependencies>编写消息发送类
// 创建连接ConnectionFactoryfactory=newConnectionFactory();factory.setHost("192.168.174.135");factory.setPort(5672);factory.setUsername("root");factory.setPassword("root");factory.setVirtualHost("/");try(Connectionconnection=factory.newConnection();Channelchannel=connection.createChannel()){/* * 定义队列 * 参数1:队列名 * 参数2:是否持久化 * 参数3:是否排他(当有一个消费者监听时,是否还可以让其他消费者监听) * 参数4:是否自动删除(如果没有任何消费者监听这个队列,是否要删除队列) * 参数5:属性,填null即可 * */channel.queueDeclare("myQueue",true,false,false,null);/* * 定义交换机 * 参数1:交换机名字 * 参数2:交换机类型 * 参数3:是否持久化 * */channel.exchangeDeclare("myExchange","direct",true);/* * 绑定队列 * 参数1:队列名字 * 参数2:交换机名字 * 参数3:routing-key(路由键) * */channel.queueBind("myQueue","myExchange","myKey");Stringmessage="hello mq";/* * 发送消息 * 参数1:交换机名称 * 参数2:路由键 * 参数3:属性,null即可 * 参数4:消息内容 * */channel.basicPublish("myExchange","myKey",null,message.getBytes(StandardCharsets.UTF_8));}catch(IOException|TimeoutExceptione){e.printStackTrace();}编写消息接收类
// 创建连接ConnectionFactoryfactory=newConnectionFactory();factory.setHost("192.168.174.135");factory.setPort(5672);factory.setUsername("root");factory.setPassword("root");factory.setVirtualHost("/");Connectionconnection=null;Channelchannel=null;try{connection=factory.newConnection();channel=connection.createChannel();channel.queueDeclare("myQueue",true,false,false,null);channel.exchangeDeclare("myExchange","direct",true);channel.queueBind("myQueue","myExchange","myKey");/* * 监听接收消息 * 参数1:队列名 * 参数2:是否自动确认 * 参数3:回调函数 */channel.basicConsume("myQueue",true,newDefaultConsumer(channel){/** * 接收消息的回调函数 * @param consumerTag 消费者编号 * @param envelope 消息的基础属性 * @param properties 基础消息的属性 * @param body 消息内容 */@OverridepublicvoidhandleDelivery(StringconsumerTag,Envelopeenvelope,AMQP.BasicPropertiesproperties,byte[]body)throwsIOException{System.out.println(newString(body,StandardCharsets.UTF_8));}});}catch(IOException|TimeoutExceptione){e.printStackTrace();}2 事务消息
事务消息与数据库的事务类似,只是MQ中的消息是要保证消息是否会全部发送成功,防止丢失消息的一种策略。
RabbitMQ有两种方式来解决这个问题:
- 通过AMQP提供的事务机制实现;
- 使用发送者确认模式实现;
事务的实现主要是对信道(Channel)的设置,主要的方法有三个:
- channel.txSelect()声明启动事务模式;
- channel.txCommint()提交事务;
- channel.txRollback()回滚事务;
3 发送者确认模式
Confirm发送方确认模式使用和事务类似,也是通过设置Channel进行发送方确认的,最终达到确保所有的消息全部发送成功
channel.confirmsSelect();// 开启发送者确认模式channel.waitForConfirms(5000L);// 确认是否发送成功waitForConfirms方法会判定在一定时间内,是否成功发送消息,如果成功返回 true,false则失败。如果抛出中断异常,那么不确定有没有发送成功,需要补发信息。
channel.addConfirmListener(newConfirmListener(){//消息确认收到后回调的方法publicvoidhandleAck(longl,booleanb)throwsIOException{System.out.println("收到消息 编号:"+l+" 是否批量:"+b);}//消息确认没有收到后的回调方法publicvoidhandleNack(longl,booleanb)throwsIOException{System.out.println("没有收到消息 编号:"+l+" 是否批量:"+b);}});addConfirmListener是异步确认,他的参数,需要定义接收成功和失败的回调函数。
4 消费者确认模式
消费者在声明队列时,可以指定 noAck 参数,当 noAck=false 时,RabbitMQ会等待消费者显式发回 ack 信号后才从内存(和磁盘,如果是持久化消息的话)中移去消息。否则,RabbitMQ会在队列中消息被消费后立即删除它。
在Consumer中Confirm模式中分为手动确认和自动确认。 手动确认主要并使用以下方法:
- basicAck(): 用于肯定确认,multiple参数用于多个消息确认。
- basicRecover():是路由不成功的消息可以使用recovery重新发送到队列中。
- basicReject():是接收端告诉服务器这个消息我拒绝接收,不处理,可以设置是否放回到队列中还是丢掉, 而且只能一次拒绝一个消息,官网中有明确说明不能批量拒绝消息,为解决批量拒绝消息才有了 basicNack。
- basicNack():可以一次拒绝N条消息,客户端可以设置basicNack方法的multiple参数为true。
channel.basicConsume(queueName,false,newDefaultConsumer(channel){publicvoidhandleDelivery(StringconsumerTag,Envelopeenvelope,AMQP.BasicPropertiesproperties,byte[]body)throwsIOException{Channelc=this.getChannel();try{System.out.println("-----准备处理消息-----");Stringmessage=newString(body);System.out.println("Receive--"+message);//获取消息的编号longmsgTag=envelope.getDeliveryTag();//手动确认消息,需要在所有的操作全部完成后将消息从队列中移除,//参数 1 为取消确认的消息编号//参数 2 为是否批量确认true表示批量确认消息,会自动移除小于等于当 前消息编号的所有消息c.basicAck(msgTag,true);}catch(Exceptione){//将消息重新放回队列,如果消息处理出现了异常则将消息从新放回队列中,尝试再次处理消息c.basicRecover();}}});四、SpringBoot集成RabbitMQ
1 配置
依赖
<dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-amqp</artifactId></dependency>配置文件
spring: rabbitmq: host: localhost port: 5672 username: root password: root配置类
@ConfigurationpublicclassRabbitCollectConfig{@BeanpublicQueuequeue(){returnnewQueue("bootQueue",true,false,false,null);}@BeanpublicDirectExchangedirectExchange(){returnnewDirectExchange("bootExchange",true,false);}@BeanpublicBindingbinding(Queuequeue,Exchangeexchange){returnnewBinding("bootQueue",Binding.DestinationType.QUEUE,exchange.getName(),"bootExchange",null);}}2 生产者
@AutowiredprivateAmqpTemplateamqpTemplate;@Testvoidsend(){amqpTemplate.convertAndSend("bootExchange","bootKey","test");}3 消费者
直接获取
@AutowiredprivateAmqpTemplateamqpTemplate;@Testvoidreceive(){amqpTemplate.receiveAndConvert("bootQueue");}或者监听
@ServicepublicclassMessageService{@RabbitListenerpublicvoidreceiveMessage(Stringmessage){System.out.println(message);}}五、使用 Canal 框架同步数据
添加依赖
<dependency><groupId>top.javatool</groupId><artifactId>canal-spring-boot-starter</artifactId><version>1.2.1-RELEASE</version></dependency>配置文件
canal:server:Canal服务部署的地址:11111destination:exampleuser-name:canalpassword:Canal_2020logging:level:root:infotop:javatool:canal:client:client:AbstractCanalClient:error添加 handler
@Slf4j@Component@CanalTable(value="t_order_info")publicclassOrderaInfoHandlerimplementsEntryHandler<OrderInfo>{@AutowiredprivateStringRedisTemplateredisTemplate;@Overridepublicvoidinsert(OrderInfoorderInfo){log.info("当有数据插入的时候会触发这个方法");}@Overridepublicvoidupdate(OrderInfobefore,OrderInfoafter){log.info("当有数据更新的时候会触发这个方法");}@Overridepublicvoiddelete(OrderInfoorderInfo){log.info("当有数据删除的时候会触发这个方法");}}编写实体类
OrderInfo,当指定的表修改之后,即可触发方法,可以发送MQ、缓存、同步其他中间件等。