ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SpringBoot整合RabbitMQ实战:交换机选型、消息可靠性与生产排障

SpringBoot整合RabbitMQ实战:交换机选型、消息可靠性与生产排障 SpringBoot整合RabbitMQ实现接发消息做后端这些年消息队列几乎是每个项目都躲不开的东西。SpringBoot整合RabbitMQ可以算是这个组合里最经典的方案之一脚手架一拉依赖就能跑但很多人的项目止步于把消息发出去、把消息收进来一旦涉及到交换机选型、消息可靠性、虚拟主机权限这些生产级问题就开始懵了。这篇文章把我从环境搭建到代码落地、再到线上排障的完整经验整理出来适合刚接触RabbitMQ的SpringBoot开发者也适合已经跑通Demo但想在项目里用得踏实一点的同学。先说说为什么选RabbitMQ而不是Kafka或RocketMQ。如果你的业务是异步任务、流量削峰、系统解耦、延迟消息这类场景RabbitMQ的延迟低、路由灵活、生态完善上手成本也低Kafka强在吞吐量和大数据管道RocketMQ强在事务消息和顺序消息的可靠性。大部分常规业务系统用RabbitMQ完全够了。这个选型问题我放到第四部分展开讲先把基础链路跑通。1. 先把环境跑起来RabbitMQ安装与后台管理里最容易被卡住的环节1.1 本地开发环境怎么装最省心先说Windows环境下怎么装RabbitMQ。RabbitMQ是Erlang语言写的所以装之前必须先装Erlang而且版本要对上。很多人栽跟头的第一站就是这里——装了一个太新的Erlang结果RabbitMQ启动直接失败日志里全是init相关的报错。我建议直接用RabbitMQ官方维护的版本对照表来选比如RabbitMQ 3.13.x对应Erlang 26.xRabbitMQ 4.0.x对应Erlang 27.x。装的时候注意两点Erlang安装路径不要带空格装到C:\erl这种简单路径避免后续环境变量解析出问题。装完Erlang后手动确认一下erl -version能正常输出再继续装RabbitMQ。安装包直接去GitHub的releases页面下载对应的Windows安装程序下载慢的话可以找国内镜像站。Mac用户直接brew install rabbitmq一条命令搞定管理员权限都帮你配置好了比Windows省心得多。最省心的方式是Docker。本地只要装了Docker Desktop一条命令就跑起来了docker run -d --name rabbitmq \ -p 5672:5672 -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:4.0-management这里解释一下端口5672是AMQP协议端口服务端和客户端通信走这个15672是Web管理界面端口浏览器访问http://localhost:15672看到的那个控制台。management后缀的镜像自带管理插件不带的版本进去还要手动rabbitmq-plugins enable rabbitmq_management多一步操作没必要。1.2 管理界面能打开但admin账号不能创建虚拟主机怎么办这个问题在相关热词里出现了我猜你们肯定有人遇到Docker部署完RabbitMQ之后浏览器能打开管理界面用admin账号也能登录但点进Virtual Hosts创建新虚拟主机它就提示没有权限或者按钮直接是灰的。先说结论镜像通过RABBITMQ_DEFAULT_USER和RABBITMQ_DEFAULT_PASS创建的这个admin用户默认只被赋予了根虚拟主机/的权限并没有管理员级别的全部权限。管理界面里很多管理操作——比如创建虚拟主机、管理其他用户、查看全局统计——都需要用户拥有administrator标签才能做。排查和修复步骤# 进入容器 docker exec -it rabbitmq bash # 查看admin用户当前权限和标签 rabbitmqctl list_users # 给admin用户打上administrator标签 rabbitmqctl set_user_tags admin administrator # 查看默认虚拟主机权限 rabbitmqctl list_permissions # 给admin配置根虚拟主机权限如果缺失 rabbitmqctl set_permissions -p / admin .* .* .*设置完之后回管理界面刷新创建虚拟主机的入口就恢复了。这里也解释一下set_permissions -p后面那个.* .* .*分别代表配置权限、写权限、读权限全给.*就是允许所有操作。生产环境按需收窄开发环境图省事全放开。还有一种情况你用的是rabbitmqctl add_user创建的用户Web管理界面却报不能连接到服务器。这种问题多半不是用户本身的问题而是节点名解析异常。RabbitMQ的节点在启动时会以主机名生成一个标识如果你改过机器名或者hosts文件里没有对应映射CLI连节点就会出现unable to connect to node这类错误。解决方法是把主机名写进hosts比如127.0.0.1 你的主机名然后重启RabbitMQ服务。1.3 Erlang版本太高或太低RabbitMQ启动失败再补充一个启动失败的常见原因RabbitMQ和Erlang版本不兼容。启动失败时的现象是服务起不来Windows事件查看器里能看到Erlang崩溃记录Docker容器则是反复重启。如果真是版本问题日志里一般会提示RabbitMQ is configured to use more than X GB memory或者直接出现The Erlang cookie相关的报错后者多半是cookie不一致。快速判断方法如果是Windows打开命令行执行erl -version和RabbitMQ版本对照表比对。如果是Docker检查一下镜像标签和对应Erlang版本的匹配关系官方镜像一般不会配错你只要别用过旧的Erlang镜像去跑新RabbitMQ就行。这里再提一个我踩过的坑Windows上如果之前装过Erlang后来升级RabbitMQ时忘了同步升级Erlang最容易出现这种问题。升级RabbitMQ之前先确认版本对照表。2. SpringBoot工程搭建依赖、配置与序列化风格的取舍2.1 依赖引入和关键配置项SpringBoot整合RabbitMQ的接入成本非常低核心依赖就一个dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency版本不用指定跟着SpringBoot的父依赖走。我是用SpringBoot 2.7.x和3.x都验证过这段代码是可以跑的Spring Boot 3.x注意JDK要17以上。然后在application.yml里写连接配置spring: rabbitmq: host: localhost port: 5672 username: admin password: admin123 virtual-host: / publisher-confirm-type: correlated publisher-returns: true listener: simple: acknowledge-mode: manual prefetch: 10 concurrency: 3 max-concurrency: 10这里的配置项一个一个说清楚virtual-host指定连接哪个虚拟主机。RabbitMQ的虚拟主机是资源隔离单元不同业务可以创建不同的虚拟主机互不干扰。开发环境用/就够了生产环境建议按业务拆分比如/order-service、/payment-service避免消息互相串。publisher-confirm-type: correlated开启发布确认。这条配置保证消息从生产者发到Broker时Broker会回调确认——这是解决消息丢失的第一道保险后面第五部分还会细讲。listener.simple.acknowledge-mode: manual消费者手动ACK。默认是自动ACK消息一被接收到就确认如果消费逻辑没处理完就抛异常消息会丢。手动ACK的意思是你在代码里确认这条消息处理成功了才告诉Broker可以删了。prefetch: 10每个消费者一次性预取的消息条数。这个参数直接影响消费吞吐量设太大会导致消息分不到其他消费者设太小浪费网络往返。concurrency和max-concurrency消费者线程数量。还记得配置里那个listener.simple吗它就是控制RabbitListener注解监听器线程池的。2.2 为什么要把默认的消息转换器换成JSONSpringBoot的AMQP starter默认带的SimpleMessageConverter用的是Java序列化也就是消息在传输前会被序列化成Java二进制格式。这个方案有两个硬伤Java序列化体积大、效率低跨语言基本没法用。RabbitMQ管理界面看到的消息内容是一堆乱码排查问题非常痛苦。所以实际项目里我第一步一定是把消息转换器替换成Jackson2JsonMessageConverter让消息在管道里以JSON文本的形式流转。操作起来很简单在配置类里声明一个MessageConverter的BeanBean public MessageConverter jacksonMessageConverter() { return new Jackson2JsonMessageConverter(); }SpringBoot会自动检测到这个Bean并注入到RabbitTemplate和监听器容器里。配置了JSON转换器之后无论发送还是接收消息体都会统一走JSON序列化/反序列化。管理界面上点开一条消息能直接看到可读的JSON内容排错的时候省下无数脑细胞。还有个细节如果你接收消息用的DTO类JSON反序列化时要注意类必须有默认构造函数或者用JsonProperty标注字段名。我遇到过同事把字段名从messageId改成msgId后消费方没同步改结果反序列化出来一堆null——问题就出在Jackson的字段映射上。2.3 CachingConnectionFactory的隐性能力配置里我们不需要显式声明连接工厂SpringBoot自动配置的CachingConnectionFactory已经帮我们创建好了。它的核心能力是缓存——它在RabbitTemplate、RabbitListener消费者之间共享一个物理连接但这个共享不是简单复用一个Connection而是维护了一组Channel。RabbitMQ的信道模型是这样的建立一条TCP长连接Connection在这条连接上开很多轻量的Channel做消息收发Channel复用远比每次都建TCP连接高效得多。你可以在代码里通过自定义ConnectionFactory来调整缓存参数Bean public CachingConnectionFactory connectionFactory(ConnectionFactoryConfigurer configurer) { CachingConnectionFactory factory new CachingConnectionFactory(); factory.setHost(localhost); factory.setPort(5672); factory.setUsername(admin); factory.setPassword(admin123); factory.setVirtualHost(/); factory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED); factory.setChannelCacheSize(25); return factory; }channelCacheSize默认是25如果并发量比较高可以适当调大。这里提醒一句话开发的时候发现一条消息收得慢先别急着调并发和prefetch先看看是不是连接工厂的缓存设置不合适再去看代码逻辑。3. 接发消息的核心代码从最简单的Hello World到带交换机路由的完整实现3.1 最基础的队列收发Hello World版本先写一个最简单的版本让消息能发出去、能收进来。这种方式不需要声明交换机消息直接发到队列里本质上用的是RabbitMQ默认的直连交换机默认交换机绑定每个队列路由键就是队列名。// 发送端 Component public class MessageSender { Autowired private RabbitTemplate rabbitTemplate; public void sendSimple(String message) { rabbitTemplate.convertAndSend(hello.queue, message); } }// 接收端 Component public class MessageReceiver { RabbitListener(queues hello.queue) public void onMessage(String message) { System.out.println(收到消息: message); } }这里有个隐含约定convertAndSend(hello.queue, message)实际上发给默认交换机路由键是hello.queue然后默认交换机会精确匹配同名的队列hello.queue。所以如果队列不存在消息会直接丢失这一点生产环境千万不要这么用。上面的写法只适合本地练手理解一下消息怎么流转就行。如果队列不存在接收端会报Reply queue not found这种错。所以实际项目中队列、交换机、绑定关系都是在配置类里预先声明的没有隐式创建这一说。3.2 完整版声明队列、交换机、绑定关系生产环境里我习惯把队列和交换机显式声明在配置类中Configuration public class RabbitConfig { public static final String EXCHANGE_NAME order.exchange; public static final String QUEUE_NAME order.queue; public static final String ROUTING_KEY order.create; Bean public DirectExchange orderExchange() { return new DirectExchange(EXCHANGE_NAME, true, false); } Bean public Queue orderQueue() { return QueueBuilder.durable(QUEUE_NAME) .deadLetterExchange(EXCHANGE_NAME.concat(.dlx)) .deadLetterRoutingKey(order.dead) .build(); } Bean public Binding orderBinding() { return BindingBuilder.bind(orderQueue()) .to(orderExchange()) .with(ROUTING_KEY); } }这里面有几个细节new DirectExchange(EXCHANGE_NAME, true, false)的三个参数分别是名称、是否持久化、是否自动删除。生产环境交换机持久化设为true自动删除设为false。为什么持久化这么重要因为如果交换机不是持久化的Broker重启后交换机就没了队列和交换机的绑定关系也就断了消息就发不进去了。QueueBuilder.durable(QUEUE_NAME)表示队列持久化。队列持久化同样是为了应对RabbitMQ重启不持久化的队列重启就消失。deadLetterExchange和deadLetterRoutingKey是死信配置消息被消费者拒绝、或者TTL过期、或者队列长度溢出时会被转投到这个死信交换机上方便做失败兜底。这个后面第五部分详细展开。3.3 发送端RabbitTemplate的用法和确认回调声明好队列和交换机之后发送端就可以用路由键来指定消息去向Service public class OrderMessageSender { Autowired private RabbitTemplate rabbitTemplate; public void sendOrderMessage(OrderDTO order) { CorrelationData correlationData new CorrelationData(UUID.randomUUID().toString()); rabbitTemplate.convertAndSend( RabbitConfig.EXCHANGE_NAME, RabbitConfig.ROUTING_KEY, order, correlationData ); } }CorrelationData是发布确认的关键。publisher-confirm-type: correlated开启模式下RabbitTemplate会在Broker确认消息之后回调CorrelationData里的ConfirmCallback。写法如下PostConstruct public void init() { rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (ack) { log.info(消息确认成功: {}, correlationData.getId()); } else { log.error(消息确认失败: {}原因: {}, correlationData.getId(), cause); } }); rabbitTemplate.setReturnsCallback(returned - { log.error(消息路由失败: {} - {}, returned.getExchange(), returned.getRoutingKey()); }); }这里要区分两个概念ConfirmCallbackBroker收到了消息返回确认ack。代表消息到了交换机。但这不保证消息一定进了队列如果路由键配错了消息会在交换机里迷路。ReturnsCallback消息从交换机路由不到任何队列时RabbitMQ会把消息退回给生产者并回调这个接口。所以完整的风控是Confirm确认到交换机Returns确认路由到队列两边都正常消息才算真正安全。3.4 接收端RabbitListener的几种姿势接收端最常见的写法Component public class OrderMessageConsumer { RabbitListener( queues RabbitConfig.QUEUE_NAME, containerFactory rabbitListenerContainerFactory ) public void onOrderMessage(OrderDTO order, Channel channel, Message message) throws IOException { long deliveryTag message.getMessageProperties().getDeliveryTag(); try { // 业务处理比如写库、调接口、更新状态 process(order); // 处理成功后手动确认 channel.basicAck(deliveryTag, false); } catch (Exception e) { // 业务失败否定确认但不重新入队等死信处理 channel.basicNack(deliveryTag, false, false); } } }这里有几个关键点方法签名可以从String message升级为OrderDTO order这个升级非常推荐因为前面配置了Jackson消息转换器框架会自动把JSON反序列化成你指定的DTO类型。这算SpringBoot为开发者省的最大的心。Channel channel这个参数不是必须的但如果你配置了手动ACK就必须接收它因为basicAck和basicNack都是通过Channel调用的。deliveryTag是消息投递标签Broker用它来定位是哪条消息。调basicAck时传false表示只确认这一条不批量确认。basicNack(deliveryTag, false, false)三个参数分别表示投递标签、是否批量、是否重新入队。我这里传的是false, false即不批量拒绝、不重新入队。为什么不重新入队因为如果消息一直消费失败重新入队会造成无限循环。正确做法是让它转死信队列后面单独处理。如果不想每次手写basicAck也可以把acknowledge-mode设为auto让Spring容器自动确认。但自动确认没法精细控制失败重投而且异常处理的语义比较隐蔽我强烈建议生产环境用手动ACK。3.5 消费者异常重试策略手动ACK模式下消息处理抛异常会走你catch里的basicNack逻辑。但你想过这个场景吗代码刚升级消费者突然连不上数据库每条消息进来都会抛异常如果你每次都basicNack并且不重新入队消息全跑到死信队列里了。这个情况下更合理的做法是让临时性故障导致的失败消息重新排队等数据库恢复了再处理。SpringBoot提供了一个优雅的重试机制在application.yml里配置spring: rabbitmq: listener: simple: retry: enabled: true max-attempts: 3 initial-interval: 1000 multiplier: 2 max-interval: 10000开启重试后消费者方法抛出异常时消息不会立刻走basicNack而是按延迟间隔重试。initial-interval: 1000表示第一次重试等1秒multiplier: 2表示每次间隔翻倍最多重试3次。重试次数耗尽之后Spring会把消息标记为失败。如果还配置了defaultRequeueRejected: false消息会继续进入死信队列。这个配置组合是生产环境的黄金搭档临时故障自动重试重试无望转入死信由专门的补偿程序去处理。4. 四种交换机不是随便选的场景、绑定关系与配置实例4.1 从RabbitMQ和Kafka的选型对比看交换机的价值网上讨论得最多的通常是RabbitMQ、Kafka、RocketMQ怎么选。我的经验是先看业务模型再看技术需求。RabbitMQ的核心竞争力是灵活的路由模型——发布/订阅、按路由键匹配、按主题匹配一套交换机机制可以覆盖很多消息分发模式Kafka更强调高吞吐日志流数据管道、实时计算场景往往依赖Kafka的分区和消费者组机制RocketMQ在事务消息、延迟消息方面做得更完善。如果只是把消息从A发给B这种简单场景Kafka其实并没有优势反而增加了学习成本和运维负担。大多数业务系统里RabbitMQ就是最优解。在RabbitMQ里交换机才决定了消息的分发方式。下面把这四种交换机讲透。4.2 Direct Exchange精确匹配直接交换机在绑定队列时需要指定一个路由键发送消息时也要指定路由键两者完全匹配消息才会进队列。适用场景订单处理、支付回调、业务类型明确的点对点通信。比如订单创建、订单支付、订单取消各绑定自己的路由键发送端用不同路由键消息就精确分流到对应消费者。Bean public Binding orderCreateBinding() { return BindingBuilder.bind(orderCreateQueue()) .to(orderDirectExchange()) .with(order.create); }Direct是默认的交换机类型也是初学者最好理解的我建议没有特殊需求之前先用它。4.3 Fanout Exchange广播扇形交换机不关心路由键它会把消息广播到所有绑定的队列上。路由键传什么都无所谓绑定多少队列消息就会复制多少份。适用场景全局通知、库存更新、缓存刷新。比如商品价格改了需要同步通知搜索服务、推荐服务、详情页缓存服务用fanout一条消息全搞定。Bean public FanoutExchange cacheFanoutExchange() { return new FanoutExchange(cache.refresh.exchange, true, false); } Bean public Binding searchCacheBinding() { return BindingBuilder.bind(searchCacheQueue()) .to(cacheFanoutExchange()); } Bean public Binding recommendCacheBinding() { return BindingBuilder.bind(recommendCacheQueue()) .to(cacheFanoutExchange()); }4.4 Topic Exchange模糊匹配主题交换机是按主题模式匹配的。路由键由点号分隔绑定关系中的路由键支持两种通配符*匹配一个单词。#匹配零个或多个单词。举例order.*能匹配order.create、order.payorder.#还能匹配order.pay.success。适用场景按业务维度灵活订阅。比如运维平台log.info.*订阅所有info日志log.error.*订阅所有error日志#.system订阅系统级所有事件组合起来非常灵活。4.5 Headers Exchange按Header匹配头部交换机不按路由键匹配而是按消息Header头里的键值对匹配。这个交换机用到的场景极少消息头匹配不如路由键直观性能和可读性都一般建议除非是遗留系统、协议兼容等问题不要选它。排个优先级能用Direct就用Direct多目标广播用Fanout需要模糊匹配用TopicHeaders除非没得选否则放弃。4.6 关于Quorum Queue的补充RabbitMQ 4.0开始官方把Quorum Queue的定位提到了很高的位置。它是基于Raft协议实现的队列类型比经典队列Classic Queue更强的一致性保证。生产环境里我建议优先考虑Quorum队列尤其是在需要高可用、防数据丢失的金融、订单类业务里。声明Quorum队列的写法Bean public Queue orderQuorumQueue() { return QueueBuilder.durable(QUEUE_NAME) .quorum() .build(); }Quorum队列注意三点不支持事务消息但发布确认是天然支持的。消息只能由主副本处理不像Kafka那样多副本分摊读流量所以不要指望它解决读扩展。和Classic队列混用时交换机绑定关系都要显式声明避免隐式绑定行为不一致。5. 消息可靠性从三个环节保住消息不丢5.1 消息从生产到消费的完整链路哪里可能丢我在前面的配置里做了一堆可靠性的铺垫发布确认、持久化队列、手动ACK、死信队列。现在把这三个环节串起来讲清楚。一条消息从生产者到消费者完整经过三段第一段生产者把消息发到RabbitMQ Broker交换机/队列。这一段丢了怎么办靠发布确认。第二段消息在队列中存储等待消费者消费。这一段丢了怎么办靠队列和消息持久化。第三段消费者拿到消息后处理。这一段丢了怎么办靠消费者手动ACK。任何一个环节出了问题消息就可能悄悄消失。我一个一个来说。5.2 发布确认确保消息到Broker前面已经配置了publisher-confirm-type: correlated和setConfirmCallback。这是生产者侧的第一道保险。确认回调里如果收到失败ackfalse可以考虑把消息存进本地一张表启动一个定时任务重发。这是非常实用的生产模式本地消息表定期补发性价比极高不需要引入复杂的分布式事务框架。5.3 队列持久化和消息持久化确保消息不随Broker重启消失队列持久化durable在前面的配置里已经做了但注意一点队列持久化不意味着消息持久化。只有在发送消息时把消息的deliveryMode设置为PERSISTENT消息才会被持久化到磁盘。如果用SpringBoot的RabbitTemplate.convertAndSend发送消息默认消息就是可持久化的。但如果你手动构造MessageProperties就要注意设置MessageProperties props new MessageProperties(); props.setDeliveryMode(MessageDeliveryMode.PERSISTENT);不设置的话消息只存在内存里一旦RabbitMQ宕机重启内存消息全部丢失。这个坑比较隐蔽因为开发环境机器一般不重启不容易暴露。5.4 消费端手动ACK和幂等性确保消息不重复消费手动ACK保证了消息处理成功才确认但分布式系统里处理成功和确认成功之间裂了一道缝消息处理完了还没来得及ACK消费者进程崩了这条消息会被重新投递给其他消费者。于是业务被重复执行了。这就是重复消费问题的根源。解决重复消费的核心不是靠MQMQ只能保证at-least-once也就是至少一次投递没法保证恰好一次而是靠消费端幂等。我在项目里最常用的方案是给业务加上唯一下标幂等判断。比如订单消息里带一个orderMessageId消费时先查Redis里这个ID有没有处理过String messageId order.getMessageId(); Boolean firstConsumed stringRedisTemplate.opsForValue() .setIfAbsent(order:msg: messageId, 1, Duration.ofDays(1)); if (firstConsumed null || !firstConsumed) { log.warn(重复消息跳过处理: {}, messageId); channel.basicAck(deliveryTag, false); return; }setIfAbsent这个操作是原子性的成功返回true说明是第一次返回false说明已经处理过了。主键唯一约束也是类似思路先查是否存在再处理存在就直接ACK跳过。幂等这一层必须做不做等于给生产环境埋雷。5.5 死信队列给处理失败的消息一个收容所死信Dead Letter这个概念前面已经出现好几次了现在完整说明一下。消息进入死信队列有三种情况消费者调用basicNack或basicReject且requeuefalse。消息设置了TTL过期时间超时未被消费。队列长度达到上限新消息无法入队。死信交换机和普通交换机没什么区别只是它专门接收死信消息。配置方式Bean public Queue orderDeadLetterQueue() { return QueueBuilder.durable(QUEUE_NAME .dlq) .build(); } Bean public DirectExchange orderDeadLetterExchange() { return new DirectExchange(EXCHANGE_NAME .dlx, true, false); } Bean public Binding orderDeadLetterBinding() { return BindingBuilder.bind(orderDeadLetterQueue()) .to(orderDeadLetterExchange()) .with(order.dead); }然后在原队列声明里加上死信参数QueueBuilder.durable(QUEUE_NAME) .deadLetterExchange(EXCHANGE_NAME .dlx) .deadLetterRoutingKey(order.dead) .build();死信队列的消费者可以做的处理包括记日志、告警让值班人员介入。把失败消息落库补偿程序定期重放。人工排查业务逻辑或代码缺陷修好之后手动搬运消息回原队列。死信队列不建议只接一个消费异常就什么都不做它的价值在于让失败消息可见、可控、可处理。6. 生产环境排障手册账号权限、连接问题与消息异常6.1 常见问题的现象、原因和处理手段把这些年在RabbitMQ上遇到的故障整理成了一张排查速查表。表里的每一行都来自真实线上事故供你参考。现象可能原因排查手段消费者收不到消息交换机路由键和绑定路由键不匹配管理界面查看队列的Binding关系确认路由键消费者收不到消息队列绑定到别的交换机了rabbitmqctl list_bindings确认绑定关系消息确认失败ackfalse交换机不存在确认交换机是否声明、是否持久化消息路由失败Returns回调触发路由键没有匹配到任何队列检查队列绑定关系或检查消息路由键拼写连接被拒绝connection refused端口和虚拟主机配置不对检查5672端口、virtual-host、用户名密码生产者超时Reply timeoutRabbitMQ负载过高或网络分区看管理界面的连接数、队列堆积量、网络指标队列里消息堆积持续上涨消费速度小于生产速度调大并发或prefetch确认消费者是否有阻塞管理界面使用rabbitmqctl创建的用户连不上节点名解析异常检查hosts映射重启RabbitMQ节点rabbitmqctl list_users能执行但web界面无法登录用户没有management或administrator标签set_user_tags补标签消息重复消费消费者处理后ACK前崩了或业务没有幂等检查是否配置手动ACK业务里增加幂等判断队列消息数量为0但业务没收到消费端抛异常后消息进了死信队列查看死信队列检查消费日志6.2 从日志到管理界面的一套排查顺序遇到RabbitMQ相关问题我最推荐的排查顺序是固定的按照这个顺序走大多数问题都能定位在几分钟内第一步看SpringBoot应用日志。消费者异常会打堆栈ConfirmCallback失败会打ackfalse和原因消息丢失如果是代码问题日志里基本都会暴露。应用日志是最先也最快的线索来源。第二步看管理界面。重点看三个地方Exchange列表里有没有对应的交换机Queue列表里有没有对应队列点击队列进去看Consumers、Messages里的Ready/Unacked状态。下面这几个指标重点说明一下Ready已经进入队列等待消费的消息数。持续增长说明消费者没有消费或消费速度跟不上。Unacked已经被消费者拉走但还没有确认的消息数。如果Unacked一直居高不下说明消费者处理很慢或者卡在某个阻塞操作上比如同步调接口超时。TotalReady Unacked的总和。这个数字和预期业务量级对比能很快判断是否消息积压。第三步如果业务日志和管理界面都看不出问题就用命令行直接查rabbitmqctl list_queues name messages_ready messages_unacknowledged rabbitmqctl list_exchanges name type rabbitmqctl list_bindings source_name destination_name routing_key rabbitmqctl list_consumers queue_name第四步抓包或开Debug日志。如果到了这步还没定位就把SpringBoot的日志级别调成Debug重点看com.rabbitmq.client.impl包的日志它能打印AMQP协议层的数据帧能判断消息是不是真的在网络上发送成功了。6.3 关于虚拟主机的一个常见误区最后单独说一句虚拟主机。很多新手把虚拟主机和队列的概念混在一起创建了一个虚拟主机然后在里面声明队列就跑不通因为连接配置里的virtual-host没有改。SpringBoot连接配置里的virtual-host必须指向你创建的那个虚拟主机。如果连接时没有指定或指错了你会看到ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN的错误。排错的时候优先确认虚拟主机名拼写和权限分配——用户名在某个虚拟主机下要有对应的权限否则连接照样拒绝。我在项目里通常这样规划开发环境一个虚拟主机测试环境另一个生产环境按业务模块再细分。这样消息在环境之间天然隔离权限控制也不容易出问题。6.4 运维上的一些经验总结队列名字要带业务前缀比如order.create.queue别用无意义的test1。这个在管理界面排障时真的能救命一眼就看出是哪个业务链路的问题。死信队列一定要配。没有死信队列的消息如果消费失败会无限重试如果requeuetrue或者直接消失如果requeuefalse两种都不是好结果。发布确认和ReturnsCallback一定要接。它们不是性能敏感操作对生产环境的排障价值却是最大的。消息DTO的字段变更要灰度发布。消费者和生产者版本不一致时旧版本消费者反序列化新字段会得到null或不报错但静默丢数据。我建议DTO增加版本号字段字段变更时保证消费者先升级兼容再升级生产者。坦率地说RabbitMQ本身的定位就是可靠、灵活、易用但能用和用得稳之间相差了一整套兜底设计。把发布确认当标配、把死信队列当默认、把幂等消费当刚需你的RabbitMQ才是真正敢上生产环境的。
返回列表