rabbitmq使用笔记

发布时间:2026/8/23 13:42:11
rabbitmq使用笔记 文章目录场景集成过程引入依赖application.properties中增加属性RabbitConfig代码DirectExchangeConfig代码订阅模式代码Service接口和实现类routingKey 路由字典表SpringBootApplication启动类MqController代码QueueListener类验证方法拓展问题是否可以通过配置的形式不消费消息单交换机和多交换机如何设计1交换机6key6queue2交换机3key6queue或者3queue6交换机不用key,6queue,或者3queuetopic模式的优点Message里面配置MessageProperties的代码listener可以直接接受Message对象么channel和queue的区别秒杀场景如何保证安全概念基本元素mq自动重连mq是http请求吗rabbitmq可以实现延迟消息吗?延迟消息有哪些应用场景?死信队列机制、底层原理报错 Failed to declare queuelistener可以配置开关吗?场景公司的数据之前是采用定时任务进行同步优点是技术比较成熟没有学习成本。但是也有不足数据无论如何会有延迟。过多定时任务也会加大服务器压力。定时任务是单线程的如果其中某条记录报错处理不好会阻断整个定时任务。用mq可以很好的解决上述问题。集成过程引入依赖pom.xml中加入dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-amqp/artifactId/dependencyapplication.properties中增加属性注以下属性不是自动注入的(不像datasource的属性),是需要手动引入的spring.rabbitmq.host192.168.10.11spring.rabbitmq.port5672spring.rabbitmq.userNamerootspring.rabbitmq.passWord1234# virtualHost起到分隔的作用因为一台机器可能供多个系统使用难到要每个服务一台机器?spring.rabbitmq.virtualHostticketRabbitConfig代码这个类主要配置连接的作用以及创建连接池/** * rabbitmq的主配置类 * 主要定义链接和虚拟机 */ConfigurationEnableRabbit// EnableRabbit 这个注解一定要加上表示开启rabbitmq服务不一定在application类上任意configuration类即可publicclassRabbitConfig{Value(${spring.rabbitmq.exchangeName})privateStringexchangeName;// 交换机其实不属于这个connection级别Value(${spring.rabbitmq.host})privateStringhost;Value(${spring.rabbitmq.port})privateIntegerport;Value(${spring.rabbitmq.userName})privateStringuserName;Value(${spring.rabbitmq.passWord})privateStringpassWord;Value(${spring.rabbitmq.virtualHost})privateStringvirtualHost;BeanpublicCachingConnectionFactoryconnectionFactory(){CachingConnectionFactoryconnectionFactorynewCachingConnectionFactory();// 这里也可以写host但是为了一致在下面写了connectionFactory.setHost(host);connectionFactory.setUsername(userName);connectionFactory.setPassword(passWord);connectionFactory.setPort(port);connectionFactory.setVirtualHost(virtualHost);returnconnectionFactory;}BeanpublicAmqpAdminamqpAdmin(){returnnewRabbitAdmin(connectionFactory());}/* 重写rabbitTemplate模板 配置上必须的参数*/// 这个待考证实测无用 不如配置Message里面的MessagePropertiesBeanpublicRabbitTemplaterabbitTemplate(){RabbitTemplaterabbitTemplatenewRabbitTemplate(connectionFactory());// 公司必须设置的参数-个性化参数MessagePropertiesmessagePropertiesnewMessageProperties();messageProperties.setAppId(ali);// 表示是从哪个系统发送的messageProperties.setContentType(text/plain);StringmessageIdUUID.randomUUID().toString();messageProperties.setMessageId(messageId);messageProperties.setContentEncoding(UTF-8);messageProperties.setType(uat.crm);messageProperties.setTimestamp(newDate());MessagePropertiesConvertermessagePropertiesConverternewDefaultMessagePropertiesConverter();messagePropertiesConverter.fromMessageProperties(messageProperties,UTF-8);rabbitTemplate.setMessagePropertiesConverter(messagePropertiesConverter);returnrabbitTemplate;}//配置消费者监听的容器BeanpublicSimpleRabbitListenerContainerFactoryrabbitListenerContainerFactory(){SimpleRabbitListenerContainerFactoryfactorynewSimpleRabbitListenerContainerFactory();factory.setConnectionFactory(connectionFactory());factory.setConcurrentConsumers(3);factory.setMaxConcurrentConsumers(10);//factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);//设置确认模式手工确认returnfactory;}}DirectExchangeConfig代码mq有多种发送规则。这里用direct模式。他和topic(订阅模式) 区别就是topic支持正则。/** * direct交换机的规则在这里配置 */ConfigurationpublicclassDirectExchangeConfig{/** * direct和topic的区别就是topic支持正则如果规则不很复杂direct就足够 * 需要多个交换器么如果规则不复杂不用多个交换器 */BeanpublicDirectExchangedirect(){DirectExchangedirectExchangenewDirectExchange(direct);returndirectExchange;}BeanpublicQueuemail(){QueuequeuenewQueue(mail);// 对列名称要和listener一致returnqueue;}BeanpublicQueueshort(){QueuequeuenewQueue(short);returnqueue;}// bingding的作用是设置规则// queue和规则是多对多关系BeanpublicBindingbindMailWithChecked(DirectExchangedirect,Queuemail){// 为什么要用with因为推送方可能推送多种如yanzhenrenzhengyanjia等不只可以推送一条消息returnBindingBuilder.bind(mail).to(direct).with(RoutingKey.MAIL);}BeanpublicBindingbindMailWithDeducted(DirectExchangedirect,Queuemail){returnBindingBuilder.bind(mail).to(direct).with(RoutingKey.SHORT);}BeanpublicBindingbindShortWithChecked(DirectExchangedirect,Queueshort){returnBindingBuilder.bind(short).to(direct).with(RoutingKey.MAIL);}BeanpublicBindingbindShortWithDeducted(DirectExchangedirect,Queueshort){returnBindingBuilder.bind(short).to(direct).with(RoutingKey.SHORT);}}订阅模式代码订阅模式其实也比较简单。可以通过路由实现分流。ConfigurationpublicclassSubscribeExchangeConfig{BeanpublicAExchangeaExchange(){FanoutExchangefanoutExchangenewFanoutExchange(aExchange);returnfanoutExchange;}BeanpublicQueuemail(){QueuequeuenewQueue(mail);// 对列名称要和listener一致returnqueue;}BeanpublicQueueshort(){QueuequeuenewQueue(short);returnqueue;}// bingding的作用是设置规则// queue和规则是多对多关系BeanpublicBindingbindMailWithAExchange(FanoutExchangeaExchange,Queuemail){returnBindingBuilder.bind(mail).to(aExchange);}BeanpublicBindingbindShortWithAExchange(FanoutExchangeaExchange,Queueshort){returnBindingBuilder.bind(short).to(aExchange);}}Service接口和实现类DirectExchangeService接口publicinterfaceDirectExchangeService{publicvoidsendMessage(StringroutingKey,Stringmessage);}DirectExchangeServiceImpl实现类ServicepublicclassDirectExchangeServiceImplimplementsDirectExchangeService{privatestaticLoggerloggerLoggerFactory.getLogger(DirectExchangeServiceImpl.class);AutowiredRabbitTemplaterabbitTemplate;AutowiredprivateDirectExchangedirectExchange;OverridepublicvoidsendMessage(StringroutingKey,Stringmessage){try{logger.info(mq准备推送数据routingKey: {} message: {},RoutingKey.MAIL,JSON.toJSONString(message));rabbitTemplate.convertAndSend(directExchange.getName(),RoutingKey.MAIL,message);}catch(Exceptione){logger.error(mq推送异常routingKey: {} ,RoutingKey.MAIL,e);}}}routingKey 路由字典表/** * routingkey建议使用字典表,作用至少有2点 * 1、后续可能有多个routingkey便于维护 * 2、send和binding中用的routingkey往往不在同一类中用string容易错漏。 */publicclassRoutingKey{publicstaticStringMAILmail;// 邮件publicstaticStringSHORTshort;// 短信}SpringBootApplication启动类注EnableRabbit注解可以加在这里也可以加载configuration类上但是一定要记得加。SpringBootApplicationpublicclassOaSpringBootApplication{publicstaticvoidmain(String[]args){SpringApplication.run(OaSpringBootApplication.class,args);}}MqController代码用来发送请求惭愧一直都不太会mock的测试方法就用请求把。ControllerRequestMapping(mq)classMqController{AutowiredprivateDirectExchangeServicedirectExchangeService;ResponseBodyRequestMapping(/send)publicStringsend(){Stringmessage1234已发邮件;directExchangeService.sendMessage(RoutingKey.MAIL,message);returnsuccess;}ResponseBodyRequestMapping(/send2)publicStringsend2(){Stringmessage666已发短信;directExchangeService.sendMessage(RoutingKey.SHORT,message);returnsuccess;}}QueueListener类listener监听相当于消费者。本例子集成在一个项目中了实际listener往往单独作为一个项目。消费者端只需要RabbitConfig和QueueListener即可。/** * 监听queue队列并实现消费推送端不需要配置listener,,因为它会消费掉队列 * 推送端请记得注释掉这个类 */ComponentpublicclassQueueListener{privatestaticLoggerloggerLoggerFactory.getLogger(QueueListener.class);RabbitListener(queuesmail)publicvoidmailListener(Stringmessage){logger.info(mail队列收到消息: {},message);// TODO 更新数据库等操作}RabbitListener(queuesshort)publicvoidbusindessListener(Stringmessage){logger.info(short队列收到消息: {},message);// TODO 更新数据库等操作}}验证方法发送mqController的请求 localhost:8080/mq/send打印listener监听到的日志即表示成功。拓展问题是否可以通过配置的形式不消费消息RabbitListener 里面是有这个方法的exclusive(排除)方法。true 排除消费false 不排除消费(默认)boolean exclusive() default false;问题是如何通过配置的形式注入呢?例如Value。注解里面不能使用Value。如果是在这里面写true或false不能适配各个环境。和直接注释掉代码的方式是一个意思。单交换机和多交换机单交换机的确定是如果交换机相同routing_key也相同。 那么所有符合该规则的都可以接受到请求。细粒度不够。多交换机配合多routing_key组合方式可以可以有 m*n种细粒度完全足够。如果情况不是很多很有一种比较无脑的方式就是不用routing_key直接交换机对应各种情况也是可以满足需求的。如何设计3个主要元素为(交换机*key一般是总数)交换机routingKeyqueue例如有2个业务3个用户对比以下几种实现方法。1交换机6key6queue功能可以实现不过要在key上花功夫2交换机3key6queue或者3queue较为合理是2*3的数量6交换机不用key,6queue,或者3queue心有点大额如果有10个业务1000个用户难道要10000个交换机。6个queue和3个queue的区别就是3个queue需要传递业务类型否则message不知道如何处理如果是6个queue那就是单业务了。topic模式的优点建议使用topic模式因为兼容性强。他比driect模式支持正则使用更灵活。比fanout细粒度更高而且也可以当fanout模式来用方法就是用* 绑定那么无论什么key都会推送到queue。就相当于fanout模式。Message里面配置MessageProperties的代码publicMessagePropertiesgenerateMessageProperties(){MessagePropertiesmessagePropertiesnewMessageProperties();messageProperties.setAppId(crm);messageProperties.setContentType(text/plain);StringmessageIdUUID.randomUUID().toString();messageProperties.setMessageId(messageId);messageProperties.setContentEncoding(UTF-8);messageProperties.setType(uat.crm);messageProperties.setTimestamp(newDate());returnmessageProperties;}publicvoidtest(){// 新建Message对象的时候加入属性MessagemessagenewMessage(.getBytes(),generateMessageProperties());}listener可以直接接受Message对象么可以以下2种写法都可以RabbitListener(queuesmail)publicvoidaListener(Messagemessage){logger.info(收到消息: {},message);// TODO 更新数据库等操作}RabbitListener(queuesshortMessage)publicvoidbListener(Stringmessage){logger.info(收到消息: {},message);// TODO 更新数据库等操作}channel和queue的区别其实这2个不是一个维度。channel 由 connection 创建channel可以发送queue。他的主要作用相当于提高性能因为如果每发送一个queue就创建一个connection,太浪费资源但是我用channel就可以保持连接不变的情况下发送多个queue。秒杀场景如何保证安全例如1秒10万条数据如果无脑推消费者服务器很可能崩。1、严格限流设置channel.basicQos(100)或prefetch1 每次只推1条或少量条给消费者绝不一次性推10万条消费者内存不会暴涨。2、手动ACK设置autoAckfalse处理完一条才basicAck 防止消费者崩溃导致丢消息。崩溃时未ACK的消息会自动重新入队安全可靠。3、异步落库消费者拿到消息后先放在本地缓存/队列里批量写入数据库 不要一条一条写而是攒够1000条或等1秒一次性批量INSERT极大降低数据库压力。概念基本元素虚拟机(以及用户)通道(channel)路由(交换机)(exchange)队列(queue)从高到低基本是这么个层级。mq自动重连其实就是几行配置也记录下吧。connectionFactory.getRabbitConnectionFactory().setAutomaticRecoveryEnabled(true);connectionFactory.getRabbitConnectionFactory().setRequestedHeartbeat(60);//60秒connectionFactory.getRabbitConnectionFactory().setConnectionTimeout(60000);//60秒connectionFactory.getRabbitConnectionFactory().setHandshakeTimeout(10000);//10秒connectionFactory.getRabbitConnectionFactory().setShutdownTimeout(10000);//10秒connectionFactory.getRabbitConnectionFactory().setRequestedChannelMax(512);//最大设置512rabbitmq的server端也会做限制最大connectionFactory.getRabbitConnectionFactory().setNetworkRecoveryInterval(10000);//10秒mq是http请求吗不是是基于长连接的所以不是http请求。rabbitmq可以实现延迟消息吗?可以。通过死信队列来实现。延迟消息有哪些应用场景?最典型的超时未支付场景。通过延迟消息触发后续逻辑。死信队列死信队列通过设置消息的生存时间TTL和死信交换机DLX实现延迟功能。具体步骤如下核心原理消息在队列中存活一段时间后TTL到期会成为死信被自动转发到指定的死信交换机再由死信交换机路由到目标队列进行消费从而实现延迟效果 。 ‌实现步骤‌1、声明死信交换机和队列‌创建死信交换机如dlx.direct和死信队列如dlx.queue并绑定 。 ‌2、‌配置普通队列的死信属性‌在普通队列如simple.queue声明时通过dead-letter-exchange指定死信交换机dead-letter-routing-key指定路由键 。 ‌3、‌设置消息或队列的TTL‌发送消息时设置expiration字段如10秒或在队列声明时设置x-message-ttl属性 。 ‌‌4、消息流转流程‌消息在普通队列中存活TTL时间后自动成为死信并被转发到死信交换机最终由死信队列的消费者处理 。 ‌代码示例队列声明时设置ttl。BeanpublicQueuetestQueue(){returnQueueBuilder.durable(TEST_QUEUE).deadLetterExchange(DEAD_EXCHANGE).deadLetterRoutingKey(DEAD_ROUTING_KEY).ttl(10000)// 10秒过期.build();}消息发送时设置ttl。rabbitTemplate.convertAndSend(simple.direct,dead,Dead letter test,msg-{msg.getMessageProperties().setExpiration(10000);// 10秒后过期returnmsg;});机制、底层原理spring在启动时在后台建立了与mq服务器的TCP长连接。这个tcp长连接是一直存在的无论是否有RabbitListener。tcp长连接的作用是建立ConnectionFactory和RabbitTemplate(不准确但是这么描述好理解些)。在tcp的基础上还有channel(虚拟通道)的概念它的作用是所有的实际操作如订阅队列收发消息确认等都由channel来完成。channel并不依赖RabbitListener如rabbitTemplate.convertAndSend()可以直接发送消息给默认交换机而不需要监听的队列。生产者的channel和消费者的channel有所不同。生产者的channel# 如rabbitTemplate.convertAndSend()用到的它是一个临时的channel用完就会归还。消费者的channel# 如RabbitListener用到的它是一个常驻的channel。报错 Failed to declare queue详细报错Caused by: org.springframework.amqp.rabbit.listener.BlockingQueueConsumer$DeclarationException: Failed to declare queue(s):[user-fanout__user__user-notify, user-fanout__role__role-notify]很明显没有找队列。例如订阅了这两个队列测试环境配了生产环境还没配就会报这个错。因为mq或eventhub是通过心跳的方式向服务器请求看是否有数据需要返回如果找不到队列就会报如上错误。另外如下报错也是同样的问题表示找不到队列。Caused by: com.rabbitmq.client.ShutdownSignalException: channel error; protocol method: #methodchannel.close(reply-code404, reply-textNOT_FOUND - no queue user-fanout__user__user-notify in vhost PROD_user-am, class-id50, method-id10)listener可以配置开关吗?可以。实现不只一种。1、通过actuator开关。 # 略2、通过ConditionalOnProperty注解实现。ConditionalOnProperty(nameuserListenerFlag,havingValuetrue,matchIfMissingtrue)这个注解加到Listener注解到的那个方法即可。