RabbitMQ 笔记

📅 2026/8/11 13:17:52
RabbitMQ 笔记
一、RabbitMQ 概述1 消息队列消息Message是指在应用间传送的数据。消息可以非常简单比如只包含文本字符串也可以更复杂可能包含嵌入对象。不建议传递对象如果需要传递复杂数据建议传递Json。消息队列Message Queue是一种应用间的通信方式消息发送后可以立即返回由消息系统来确保消息的可靠传递。消息发布者只管把消息发布到 MQ 中而不用管谁来取消息使用者只管从 MQ 中取消息而不管是谁发布的。这样发布者和使用者都不用知道对方的存在。消息队列用于业务解耦、最终一致性、广播、错峰流控等等情况2 RabbitMQ 特点RabbitMQ 是一个由 Erlang 语言开发的 AMQP 的开源实现。AMQP Advanced Message Queue高级消息队列协议。它是应用层协议的一个开放标准为面向消息的中间件设计基于此协议的客户端与消息中间件可传递消息并不受产品、开发语言等条件的限制。RabbitMQ 最初起源于金融系统用于在分布式系统中存储转发消息在易用性、扩展性、高可用性等方面表现不俗。具体特点包括可靠性ReliabilityRabbitMQ 使用一些机制来保证可靠性如持久化、传输确认、发布确认。灵活的路由Flexible Routing在消息进入队列之前通过 Exchange 来路由消息的。对于典型的路由功能RabbitMQ 已经提供了一些内置的 Exchange 来实现。针对更复杂的路由功能可以将多个 Exchange 绑定在一起也通过插件机制实现自己的 Exchange。消息集群Clustering多个 RabbitMQ 服务器可以组成一个集群形成一个逻辑 Broker。高可用队列可以在集群中的机器上进行镜像使得在部分节点出问题的情况下队列仍然可用。多种协议Multi-protocolRabbitMQ 支持多种消息队列协议比如 STOMP、MQTT 等等。多语言客户端Many ClientsRabbitMQ 几乎支持所有常用语言比如 Java、.NET、Ruby 等等。管理界面Management UIRabbitMQ 提供了一个易用的用户界面使得用户可以监控和管理消息 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 类型转发消息是最快的。topictopic 交换器通过模式匹配分配消息的路由键属性将路由键和某个模式进行匹配此时队列需要绑定 到一个模式上。它将路由键和绑定键的字符串切分成单词这些单词之间用点隔开。它同样也会识别两个通配符符号“#”和符号“*”。#匹配0个或多个单词*匹配不多不少一个单词。三、Java RabbitMQ1 基础操作依赖dependenciesdependencygroupIdcom.rabbitmq/groupIdartifactIdamqp-client/artifactIdversion5.1.1/version/dependency/dependencies编写消息发送类// 创建连接ConnectionFactoryfactorynewConnectionFactory();factory.setHost(192.168.174.135);factory.setPort(5672);factory.setUsername(root);factory.setPassword(root);factory.setVirtualHost(/);try(Connectionconnectionfactory.newConnection();Channelchannelconnection.createChannel()){/* * 定义队列 * 参数1队列名 * 参数2是否持久化 * 参数3是否排他当有一个消费者监听时是否还可以让其他消费者监听 * 参数4是否自动删除如果没有任何消费者监听这个队列是否要删除队列 * 参数5属性填null即可 * */channel.queueDeclare(myQueue,true,false,false,null);/* * 定义交换机 * 参数1交换机名字 * 参数2交换机类型 * 参数3是否持久化 * */channel.exchangeDeclare(myExchange,direct,true);/* * 绑定队列 * 参数1队列名字 * 参数2交换机名字 * 参数3routing-key(路由键) * */channel.queueBind(myQueue,myExchange,myKey);Stringmessagehello mq;/* * 发送消息 * 参数1交换机名称 * 参数2路由键 * 参数3属性null即可 * 参数4消息内容 * */channel.basicPublish(myExchange,myKey,null,message.getBytes(StandardCharsets.UTF_8));}catch(IOException|TimeoutExceptione){e.printStackTrace();}编写消息接收类// 创建连接ConnectionFactoryfactorynewConnectionFactory();factory.setHost(192.168.174.135);factory.setPort(5672);factory.setUsername(root);factory.setPassword(root);factory.setVirtualHost(/);Connectionconnectionnull;Channelchannelnull;try{connectionfactory.newConnection();channelconnection.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方法会判定在一定时间内是否成功发送消息如果成功返回 truefalse则失败。如果抛出中断异常那么不确定有没有发送成功需要补发信息。channel.addConfirmListener(newConfirmListener(){//消息确认收到后回调的方法publicvoidhandleAck(longl,booleanb)throwsIOException{System.out.println(收到消息 编号:l 是否批量b);}//消息确认没有收到后的回调方法publicvoidhandleNack(longl,booleanb)throwsIOException{System.out.println(没有收到消息 编号:l 是否批量b);}});addConfirmListener是异步确认他的参数需要定义接收成功和失败的回调函数。4 消费者确认模式消费者在声明队列时可以指定 noAck 参数当 noAckfalse 时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{Channelcthis.getChannel();try{System.out.println(-----准备处理消息-----);StringmessagenewString(body);System.out.println(Receive--message);//获取消息的编号longmsgTagenvelope.getDeliveryTag();//手动确认消息需要在所有的操作全部完成后将消息从队列中移除//参数 1 为取消确认的消息编号//参数 2 为是否批量确认true表示批量确认消息会自动移除小于等于当 前消息编号的所有消息c.basicAck(msgTag,true);}catch(Exceptione){//将消息重新放回队列如果消息处理出现了异常则将消息从新放回队列中尝试再次处理消息c.basicRecover();}}});四、SpringBoot集成RabbitMQ1 配置依赖dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-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 框架同步数据添加依赖dependencygroupIdtop.javatool/groupIdartifactIdcanal-spring-boot-starter/artifactIdversion1.2.1-RELEASE/version/dependency配置文件canal:server:Canal服务部署的地址:11111destination:exampleuser-name:canalpassword:Canal_2020logging:level:root:infotop:javatool:canal:client:client:AbstractCanalClient:error添加 handlerSlf4jComponentCanalTable(valuet_order_info)publicclassOrderaInfoHandlerimplementsEntryHandlerOrderInfo{AutowiredprivateStringRedisTemplateredisTemplate;Overridepublicvoidinsert(OrderInfoorderInfo){log.info(当有数据插入的时候会触发这个方法);}Overridepublicvoidupdate(OrderInfobefore,OrderInfoafter){log.info(当有数据更新的时候会触发这个方法);}Overridepublicvoiddelete(OrderInfoorderInfo){log.info(当有数据删除的时候会触发这个方法);}}编写实体类OrderInfo当指定的表修改之后即可触发方法可以发送MQ、缓存、同步其他中间件等。