RabbitMQ笔记

发布时间:2026/10/11 1:17:43
RabbitMQ笔记
RabbitMQ笔记来源https://www.bilibili.com/video/BV15k4y1k7Ep/1、基础MQ基本概念RabbitMQ快速入门RabbitMQ的工作模式Spring整合RabbitMQ1.1 MQ概述MQ全称Message Queue消息队列是在消息的传输过程中保存消息的容器多用于分布式系统之间进行通信小结MQ消息队列存储消息的中间件分布式系统通信两种方式同步通信RPC框架和HTTP REST、异步通信消息队列在MQ中发送方称为生产者接收方称为消费者补充本地调用调用的函数或方法在当前程序进程内执行不需要通过网络远程调用调用的功能在另一个进程或另一台计算机上执行通常需要通过网络进行通信1.2 MQ的优势和劣势优势应用解耦异步提速削峰填谷劣势系统可用性降低系统引入外部依赖越多系统面临风险就越高。比如对于第一个例子如果MQ挂掉系统也会崩溃系统复杂度提高需要考虑有没有重复消费消息丢失消息传递顺序等等一致性问题对于场景21号系统处理完后成功但对于234号系统如果其中一个写库失败就会导致数据不一致。1.3 MQ的优势1. 应用解耦使用MQ是的应用间解耦提升容错性和可维护性2. 异步提速用户点击完下单按钮后只需等待25ms就能得到下单响应205 25ms。提升用户的体验和系统吞吐量单位时间内处理请求的数目。3. 削峰填谷使用了MQ之后限制消费消息的速度为1000这样一来高峰期产生的数据势必会被积压在MQ中高峰就被“削”掉了但是因为消息积压在高峰期过后的一段时间内消费消息的速度还是会维持在1000直到消费完积压的消息这就叫做“填谷”。使用MQ后可以提高系统稳定性。1.4 MQ的劣势系统的可用性降低系统引入的外部依赖越多系统稳定性越差一旦MQ宕机就会对业务造成影响。如何保证MQ的高可用系统复杂度提高MQ的加入大大增加了系统的复杂度以前系统间是同步的远程调调用现在是通过MQ进行异步调用如何保证消息没有被重复消费怎么处理消息丢失情况如何保证消息传递的顺序性一致性问题A系统处理完业务通过MQ给B、C、D三个系统发消息数据如果B系统、C系统处理成功D系统处理失败如何保证消息数据处理的一致性小结既然MQ有优势也有劣势那么使用MQ需要满足什么条件呢生产者不需要从消费者出获得反馈。引入消息队列之前的直接调用其接口的返回值应该是空这才让明明下层的动作还没做上层却当成动作做完了继续往后走即所谓异步成为了可能容忍数据短暂的不一致性确实是用了有效果即解耦、提速、削峰这些方面的受益超过加入MQ管理MQ这些成本。1. 5 常见的MQ产品目前业界有许多的MQ产品例如RabbitMQ、RocketMQ、ActiveMQ、Kafka、ZeroMQ、MetaMQ等也有直接使用redis充当消息队列的案例而这些消息队列产品各有侧重在实际选型时需要结合自身需求及MQ产品特征综合考虑。1.6 RabbitMQ简介AMQP即Advanced Message Queuing Protocol高级消息队列协议是一个网络协议是应用层协议的一个开放标准为面向消息的中间件设计。基于此协议的客户端与消息中间件可传递消息并不受客户端/中间件不同产品不同的开发语言等条件的限制。2006年AMQP规范发布类比HTTP。2007年Rabbit技术公司基于AMQP标准开发的RabbitMQ采用Erlang语言开发。Erlang语言由Ericson设计专门为开发高并发和分布式系统的一种语言在电信领域使用广泛。RabbitMQ基础架构如下图RabbitMQ中的相关概念Broker接受和分发消息的应用RabbitMQ Server就是Message BrokerVirtual host出于多租户和安全因素设计的把AMQP的基本组件划分到一个虚拟的分组中类似于网络中的namespace概念。当多个不同的用户使用同一个RabbitMQ Server提供的服务时可以划分出多个vhost每个用户在自己的vhost创建exchange/queue等。Connectionpublish/consumer和broker之间的TCP连接Channel如果每一次访问RabbitMQ都建立一个Connection在消息量大的时候建立TCP Connection的开销将是巨大的效率也比较低。Channel是在connection内部建立的逻辑连接如果应用程序支持多线程通常每个thread创建单独的channel进行通讯AMQP method包含了channel id帮助客户端和message broker识别channel所以channel之间是完全隔离的。Channel作为轻量级的Connection极大减少了操作系统建立TCP connection的开销。Exchangemessage到达broker的第一站根据分发规则匹配查询表中的routing key分发消息到queue中去。常用的类型有directpoint-to-pointtopicpublish-subscribeand fanoutmulticast;Queue消息最终被送到这里等待consumer取走Bindingexchange和queue之间的虚拟连接binding中可以包含routing key。Binding信息被保存到exchange中的查询表中用于message的分发依据RabbitMQ提供了6中工作模式简单模式、工作队列模式、发布订阅模式、路由模式、主题模式、RPC远程调用模式远程调用不太算MQ暂不作介绍官网对应模式介绍https://www.rabbitmq.com/getstarted.html1.7 JMSJMS即Java消息服务JavaMessage Service应用程序接口是一个Java平台中关于面向消息中间件的APIJMS是JavaEE规范中的一种类比JDBC很多消息中间件都实现了JMS规范例如ActiveMQ、RabbitMQ官方没有提供JMS的实现包但是开源社区有小结RabbitMQ是基于AMQP协议使用Erlang语言开发的一款消息队列产品RabbitMQ提供了6种工作模式我们学习5种AMQP是协议类比HTTPJMS是API规范接口类比JDBS2、 RabbitMQ的安装略3、入门程序需求使用简单模式完成消息传递步骤创建工程生产者、消费者分别添加依赖编写生产者发送消息编写消费者接受消息原生纯Java方式4、模式介绍4.0 简单模式默认的exchangeDirect类型:如果用空字符串去声明一个exchange那么系统就会使用””AMQP default”这个exchange我们创建一个queue时,默认的都会有一个和新建queue同名的routingKey绑定到这个默认的exchange上去4.1 Work Queues工作队列模式在RabbitMQ中生产者发送消息不会直接将消息投递到队列中而是先将消息投递到交换机中 在由交换机转发到具体的队列 队列再将消息以推送或者拉取方式给消费者进行消费1. 模式说明Work Queues与入门程序的简单模式相比多了一个或一些消费端多个消费端共同消费同一个队列中的消息。一个消息要么被A消费要么被B消费应用场景对于任务过重或任务较多情况使用工作队列可以提高任务处理的速度。4.2 发布订阅模式1. 模式说明在订阅模式中多了一个Exchange角色而且过程略有变化P生产者也就是要发送消息的程序但是不再发送到队列中而是发给X交换机C消费者消息的接收这会一直等待消息到来Queue消息队列接受消息、缓存消息Exchange交换机X。一方面接收生产者发送的消息。另一方面知道如何处理消息例如递交给某个特别队列、递交给所有队列、或是将消息丢弃。到底如何操作取决于Exchange的类型。Exchange有常见以下3中类型Fanout广播将消息交给所有绑定到交换机的队列Direct定向把消息交给符合指定routing key的队列Topic通配符把消息交给符合routing pattern路由模式的队列Exchange交换机只负责转发消息不具备存储消息的能力因此如果没有任何队列那么消息会丢失4.3 Routing路由模式1. 模式说明队列与交换机的绑定不能是任意绑定了而是要指定一个RoutingKey路由key;消息的发送方在向Exchange发送消息时也必须指定消息的RoutingKeyExchange不再把消息交给每一个绑定的队列而是根据消息的Routing Key进行判断只有队列的Routing Key与消息的Routing Key完全一致才会接收到消息4.4 Topic通配符模式1. 模式说明图解红色Queue绑定的是usa.#因此凡是以usa开头的routing key都会被匹配到黄色Queue绑定的是#.news因此凡是以.news结尾的routing key都会被匹配小结Topic主题模式可以实现Pub/Sub发布与订阅模式和Routing路由模式的功能只是Topic在配置routing key的时候可以使用通配符显得更加灵活。#表示一个或多个词*表示一个词多个字符需要用 “.” 连接5、SpringBoot整合RabbitMQ生产者1. 创建生产者SpringBoot工程2. 引入依赖坐标!-- https://mvnrepository.com/artifact/org.springframework.boot/spring-boot-starter-amqp --dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-amqp/artifactIdversion3.0.4/version/dependency编写properties配置配置基本信息#http://192.168.121.129:15672/ 管理界面spring.rabbitmq.host192.168.121.129spring.rabbitmq.port5672spring.rabbitmq.usernameadminspring.rabbitmq.passwordadmin spring.rabbitmq.virtual-host/lq#消息可靠投递 确认模式spring.rabbitmq.publisher-confirm-typecorrelated#消息可靠投递 回退模式spring.rabbitmq.publisher-returnstrue定义交换机队列以及绑定关系的配置类注入RabbitTemplate调用方法完成消息发送消费者创建消费者SpringBoot工程引入start依赖!-- https://mvnrepository.com/artifact/org.springframework.boot/spring-boot-starter-amqp --dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-amqp/artifactIdversion3.0.4/version/dependency编写properties配置配置基本信息spring.rabbitmq.host192.168.121.129 spring.rabbitmq.port5672 spring.rabbitmq.usernameadmin spring.rabbitmq.passwordadmin spring.rabbitmq.virtual-host/lq #消费端限流每次接收5个 spring.rabbitmq.listener.simple.prefetch5定义监听类使用RabbitListener注解完成队列监听小结SpringBoot提供了快速整合RabbitMQ的方式基本信息在properties中配置队列交换机以及绑定关系在配置类中使用Bean的方式配置生产端直接注入RabbitTemplate完成消息发送消费端直接使用RabbitListener完成消息接收6、高级RabbitMQ高级特性消息可靠性投递在使用RabbitMQ的时候作为消息发送方希望杜绝任何消息丢失或者投递失败场景。RabbitMQ为我们提供了两种模式来控制消息的可靠投递。confirm 确认模式return 回退模式RabbitMQ整个消息投递的路径为producer–rabbitmq broker–exchange–queue–consumer- 消息从producer到exchange则会返回一个confirmCallback- 消息从exchange–queue投递失败则会返回一个returnCallback我们将利用这两个callback控制消息的可靠性投递消息可靠投递之确认模式在properties中开启确认模式#消息可靠投递 确认模式 spring.rabbitmq.publisher-confirm-typecorrelated在发送消息的时候在rabbitTemplate中添加确认回调逻辑ResponseBodyRequestMapping(/sendConfirm)publicStringsendConfirm(){//设置确认模式消息从生产者到交换机rabbitTemplate.setConfirmCallback(newRabbitTemplate.ConfirmCallback(){Overridepublicvoidconfirm(CorrelationDatacorrelationData,booleanack,Stringcause){System.out.println(JSON.toJSONString(correlationData));if(ack){System.out.println(消息已经成功传入交换机);}System.out.println(cause);}});rabbitTemplate.convertAndSend(RabbitMQConfig.BASIC_EXCHANGE,basic.hello,hello,confirm);returnsuccess;}消息可靠投递之回退模式在properties中开启回退模式#消息可靠投递 回退模式 spring.rabbitmq.publisher-returnstrue在rabbitTemplate中添加回退逻辑ResponseBodyRequestMapping(/sendReturn)publicStringsendReturn(){rabbitTemplate.setMandatory(true);//设置回退模式消息从交换机到队列rabbitTemplate.setReturnsCallback(newRabbitTemplate.ReturnsCallback(){OverridepublicvoidreturnedMessage(ReturnedMessagereturned){//正常的消息发送不会执行此处的逻辑System.out.println(JSON.toJSONString(returned));}});rabbitTemplate.convertAndSend(RabbitMQConfig.BASIC_EXCHANGE,basic.hello,hello,return);returnsuccess;}Consumer ACKack指Acknowledge确认表示消费端收到消息后的确认方式。有三种确认方式- 自动确认acknowledge“none”- 手动确认acknowledge“manual”- 根据异常情况确认acknowledge“auto”这种方式使用麻烦不作讲解其中自动确认是指当消息一旦被Consumer接收到则自动确认收到并将相应message从RabbitMQ的消息缓存中移除。但是在实际业务处理中很可能消息接收到业务处理异常那么该消息就会丢失。如果设置了手动确认方式则需要在业务处理成功后调用channel.basicAck()手动签收如果出现异常则调用channel.basicNack()方法让其自动重新发送消息。Consumer Ack 小结- 在RabbitListener(queues “basic-queue”,ackMode “MANUAL”)设置ack方式none自动确认manual手动确认- 在服务端没有出现异常则调用channel.basicAck(deliveryTag,false)方式确认签收消息- 如果出现异常则在catch中调用basicNack或baxicReject拒绝消息让MQ重新发送消息消息可靠性总结持久化exchange要持久化queue要持久化message要持久化生产方确认Confirm消费方确认AckBroker高可用补充小知识可以设置消费者basicNack之后的重试次数以及时间间隔https://blog.csdn.net/Hmj050117/article/details/121589463消费端限流具体操作在application.properties中配置spring.rabbitmq.listener.simple.prefetch5消息端的确认模式一定为手动确认可以在消费者监听中配置如下RabbitListener(queuesbasic-queue,ackModeMANUAL)publicvoidreceiveMsg(Messagemsg,Channelchannel)throwsIOException{....}// basic-queue为队列名TTLTTL全称Time To Live存活时间/过期时间当消息到达存活时间后还没有被消费会被自动清除RabbitMQ可以对消息设置过期时间也可以对整个队列Queue设置过期时间。在RabbitMQConfig中设置整个队列中消息的过期时间Bean(ttlQueue)publicQueuettlQueue(){//设置队列中消息的过期时间为10sreturnQueueBuilder.durable(ttl-queue).ttl(10_000).build();}设置单个消息的过期时间:rabbitTemplate.convertAndSend(RabbitMQConfig.BASIC_EXCHANGE,ttl.hello,hello,ttl,newMessagePostProcessor(){OverridepublicMessagepostProcessMessage(Messagemessage)throwsAmqpException{//单位毫秒message.getMessageProperties().setExpiration(5000);returnnull;}});死信队列死信队列英文缩写DLXDead Letter Exchange死信交换机当消息成为Dead messge后可以被重新发送到另一个交换机这个交换机就是DLX。消息成为死信的三种情况队列消息长度达到限制消费者拒接消费消息basicNack/basicReject并且不把消息重新放入原目标队列requeuefalse原队列存在消息过期设置消息到达超时时间未被消费队列绑定死信交换机给队列设置参数Springx-dead-letter-exchange和x-dead-letter-routing-keySpringBootBean(ttlQueue1)publicQueuettlQueue1(){//设置队列中消息的过期时间为10sreturnQueueBuilder.durable(ttl-queue1).deadLetterExchange(dead-exchange).deadLetterRoutingKey(deadInt.#).maxLength(6).build();}注意死信队列与正常队列并无本质的区别只是换了一种说法而已延迟队列延迟队列即消息进入队列后不会立即被消费只有到达指定时间后才会被消费。TTL死信队列即可完成需求1. 下单后30分钟未支付取消订单回滚库存2. 新用户注册成功7天后发送短信问候实现方式1. 定时器2. 延迟队列很可惜在RabbitMQ中并未提供延迟队列功能。 但是可以使用TTL死信队列组合实现延迟队列的效果。消息可靠性保障消息幂等性处理两种方式发送消息前生产一个随机的key保存到redis中之后在处理消息的时候先查看key是否存在。不存在则不做任何处理如果存在则进行逻辑操作完成之后将key删除掉通过mysql的唯一索引发送消息前生产一个随机的key在处理消息的时候将key插入到数据库中该key对应的是mysql中的唯一索引列如果插入成功则执行逻辑操作