ARTICLE DETAIL

资讯详情

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

Spring Boot 集成 RabbitMQ:从消息队列到死信队列的完整实战指南

Spring Boot 集成 RabbitMQ:从消息队列到死信队列的完整实战指南 Spring Boot 项目里接 RabbitMQ基本是后端开发绕不开的活儿。我这两年经手过好几个 Spring Boot 中间件相关的项目从订单异步通知到日志收集再到定时任务的削峰RabbitMQ 都是主力。这篇就基于我实际开发和线上运维的经验把 Spring Boot 集成 RabbitMQ 从环境搭建到生产者消费者完整链路、再到手动确认、重试机制、死信队列这些进阶配置全部梳理一遍最后整理一批高频故障的排查思路。适合正在做 Spring Boot 项目、打算引入消息队列的开发者也适合那些已经被 RabbitMQ 各种诡异问题折磨过的同学。1. 为什么要在 Spring Boot 项目里用 RabbitMQ1.1 RabbitMQ 在项目中到底解决什么问题先说一个很现实的场景之前做一个上门烹饪预约系统用户下单后要通知厨师、发送短信、更新排班状态。如果这些操作全部同步执行高峰期一个下单接口可能耗时好几秒用户体验非常差。引入 RabbitMQ 之后下单接口只负责写入订单、把消息投递到队列剩下的短信通知、厨师分配全部异步消费接口响应时间直接从两秒多降到了两百毫秒以内。这个例子其实点出了 RabbitMQ 的核心定位它是个消息代理帮我们把“产生事件”和“处理事件”解耦。放在 Spring Boot 项目里最典型的用处有三个。第一是异步处理把不需要同步返回结果的耗时操作丢到队列里比如发邮件、写日志、生成报表。第二是流量削峰突发流量先打到 MQ后端消费者按自己的能力慢慢处理比如秒杀场景的库存扣减。第三是应用解耦生产者和消费者互不感知对方的存在服务 A 挂了不影响服务 B 继续收消息。很多初学者会问我用 Redis 的 List 或者 Stream 不也能做队列吗确实能但 Redis Stream 更适合轻量级场景RabbitMQ 胜在完整的 AMQP 协议支持、成熟的路由模型、消息确认机制、死信队列这些企业级特性。如果你需要精细的路由策略、消息不丢失、以及复杂的重试死信机制RabbitMQ 是更稳的选择。1.2 Spring Boot 集成 RabbitMQ 的选型思路Spring Boot 对 RabbitMQ 的封装非常完善我们只需要引入spring-boot-starter-amqp这个依赖加上简单的配置就能获得一个开箱即用的 RabbitTemplate 用于发送消息以及 RabbitListener 注解用于消费消息连连接管理都不用自己操心底层自动配置会帮你创建 ConnectionFactory。这里有个选型上的经验点。在 Spring Boot 项目里用 RabbitMQ我建议直接用官方 starter 而不是自己封装一套客户端。官方 starter 帮你处理好了连接池、自动重连、JSON 序列化配置这些都是生产环境必须的基础能力。如果自己写很容易在连接管理上踩坑。另一个决策点是交换机类型的选择。RabbitMQ 提供 Direct、Topic、Fanout、Headers 四种交换机我在实际项目里用最多的是 Direct 和 Topic。Direct 适合点对点的精确路由比如“订单创建”事件发给“订单处理队列”Topic 适合带通配符的灵活路由支持order.*、order.#这样的模式匹配比如日志系统里不同级别的日志路由到不同队列。Fanout 则是广播所有绑定的队列都能收到同一条消息适合需要多服务同时感知同一事件的场景。2. 环境准备RabbitMQ 安装与 Spring Boot 基础配置2.1 Docker 安装 RabbitMQ 与端口说明新项目我基本都用 Docker 部署 RabbitMQ好处是环境一致性好清理起来也干净。最常用的安装命令是这样的docker run -d --name rabbitmq \ -p 5672:5672 \ -p 15672:15672 \ -p 25672:25672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:3.12-management这里5672是 AMQP 协议端口应用程序连接就用它15672是管理界面端口浏览器访问可以查看队列、交换机、消息状态25672是集群节点间通信端口。如果你要搭建集群除了这三个还要额外暴露4369epmd 节点发现以及1883等插件端口。镜像我建议直接带management后缀否则装完还得自己进去启用管理插件麻烦。如果你需要延迟消息插件可以在容器里执行rabbitmq-plugins enable rabbitmq_delayed_message_exchange或者用包含插件的自定义镜像这点在做延迟队列时特别有用。2.2 Windows 本地安装 RabbitMQ在公司 Windows 电脑上调试时Docker Desktop 有时候因为虚拟化问题起不来这时候就得在本机装了。Windows 上安装 RabbitMQ 依赖 Erlang 环境版本对应关系在 RabbitMQ 官网有对照表装错了会启动失败。大概步骤是先下载并安装匹配的 Erlang然后安装 RabbitMQ 的 Windows 安装包装完后进入 RabbitMQ 安装目录的 sbin 文件夹执行rabbitmq-plugins enable rabbitmq_management启用管理界面最后访问http://localhost:15672默认账号guest/guest。这里有个坑默认情况下 guest 账号只允许本机访问。如果你的 Spring Boot 应用跑在远程机器上用 guest 连过来会被拒绝。解决方法是新建一个拥有权限的用户或者配置 loopback_users 列表我一般直接用 Docker 或 rabbitmqctl 创建独立用户避免动全局配置。2.3 Spring Boot 依赖与 yml 配置在 Spring Boot 项目里引入 RabbitMQ 很简单pom 里加一个依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency然后配置 application.ymlspring: rabbitmq: host: localhost port: 5672 username: admin password: admin123 virtual-host: / publisher-confirm-type: correlated publisher-returns: true listener: simple: acknowledge-mode: manual retry: enabled: true max-attempts: 3 initial-interval: 1000 multiplier: 2 max-interval: 10000注意到几个重点参数publisher-confirm-type: correlated开启发布确认消息发送到交换机后回调确认publisher-returns: true开启消息未路由到队列的回退机制acknowledge-mode: manual设置消费者手动确认。后面会专门讲这些配置的实际意义这里先埋个伏笔。3. 核心实现从生产者到消费者的完整链路3.1 交换机、队列与绑定关系的设计RabbitMQ 里消息不是直接发给队列的而是发给交换机交换机根据路由键把消息投递到绑定的队列。这个机制很多人刚学时觉得绕其实可以把它想象成快递分拣中心你寄快递时只填目的地路由键分拣中心交换机根据地址把包裹发到对应的运输线路队列上。在 Spring Boot 里我们可以用代码声明这些组件。我习惯把队列、交换机、绑定的定义写在一个配置类里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).build(); } Bean public Binding orderBinding() { return BindingBuilder.bind(orderQueue()) .to(orderExchange()) .with(ROUTING_KEY); } }这里我把交换机、队列都设为持久化durable true就算 RabbitMQ 重启这些元数据也不会丢。队列用QueueBuilder.durable()还有一个额外的好处可以链式调用来设置死信参数后面讲死信队列时再展开。3.2 生产者代码实现与 Confirm 回调生产者投递消息时最常用的工具就是RabbitTemplate。在 Spring Boot 中我们可以直接注入使用。一条消息从应用发出去其实要经历两次确认第一次是消息发送到交换机成功Broker 返回确认回调第二次是消息从交换机路由到队列成功如果不成功就触发 ReturnedMessage 回调。这两个机制是保证消息不丢的关键。Service public class OrderProducer { private final RabbitTemplate rabbitTemplate; public OrderProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate rabbitTemplate; } public void sendOrderMessage(OrderMessage message) { CorrelationData correlationData new CorrelationData(UUID.randomUUID().toString()); rabbitTemplate.convertAndSend( RabbitConfig.EXCHANGE_NAME, RabbitConfig.ROUTING_KEY, message, correlationData ); // 通过回调确认消息状态 rabbitTemplate.setConfirmCallback((correlation, ack, cause) - { if (ack) { log.info(消息发送成功, id: {}, correlation.getId()); } else { log.error(消息发送失败: {}, cause); // 这里可以记录日志并手动补偿 } }); } }注意convertAndSend会使用 Spring Boot 配置的消息转换器。如果直接用 JDK 默认的序列化消息在管理界面上看到的是乱码而且性能差。我在项目里会配置 Jackson JSON 转换器让消息以 JSON 格式传输。具体做法是定义一个Jackson2JsonMessageConverterBean然后设置到 RabbitTemplate 和 SimpleRabbitListenerContainerFactory 上。3.3 消费者代码实现与并发参数调整消费端最简单的方式是用 RabbitListener 注解直接监听队列。在手动确认模式下消费者处理完消息后必须调用basicAck告诉 Broker 这条消息处理成功否则消息会一直留在队列里。Component public class OrderConsumer { RabbitListener(queues RabbitConfig.QUEUE_NAME) public void onMessage(OrderMessage message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) { try { // 业务处理更新订单状态、发送通知 processOrder(message); // 手动确认false 表示不批量确认 channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(处理订单消息失败, e); // basicNack 的第三个参数 是否重新入队这里先设为 false配合死信队列使用 channel.basicNack(deliveryTag, false, false); } } }很多新手会忽略一个关键参数prefetchCount。这个参数决定了消费者当前可以“预取”多少条消息默认值是 250意味着一个消费者会一次性拉取大量消息到本地内存。如果你的消息处理速度跟不上或者内存有限建议把这个值调小比如 10 或者 20这样 RabbitMQ 会慢慢分发消息避免瞬时负载过高。可以在 yml 的 listener 配置里添加prefetch: 10。并发方面RabbitListener默认是单线程消费。如果队列消息量大可以配置并发消费者数spring: rabbitmq: listener: simple: concurrency: 5 max-concurrency: 10这个配置代表每个监听容器最少 5 个、最多 10 个并发消费线程。但要注意并发数不是越大越好还要看下游数据库、接口的承受能力否则消费速度快了下游资源反而被压垮这就是典型的削峰逻辑要权衡的点。4. 手动确认、重试机制与死信队列4.1 手动 ACK 为什么是必须的RabbitMQ 的消息确认机制有三种自动确认、手动确认、以及不确认。Spring Boot 默认是自动确认也就是消费者拿到消息后无论处理成功与否Broker 都认为消息已被消费。这在业务代码抛异常时会直接造成消息丢失而且是在还没记录任何日志的情况下丢的。我之所以强调手动确认是因为实际生产环境中业务处理难免出错我们大概率需要失败重试或者进入死信队列。自动确认的模式下异常消息直接当作“成功”删除了后面想排查都无从下手。手动确认则把决定权交给我们处理成功就 ack处理失败就 nack 或者 basicReject。把acknowledge-mode设为manual之后消费者方法签名必须添加 Channel 参数和 Header 注解拿到 deliveryTag因为在手动模式下确认动作必须通过 Channel 调用 basicAck 或 basicNack 完成。确认时第二个参数multiple要明白它的含义传 true 表示确认当前消息之前的所有消息传 false 只确认当前这一条。我在项目里统一用 false避免误确认未处理的消息。4.2 重试机制的两种实现方式消费失败后的重试大致有两种思路RabbitMQ 侧的重试和 Spring 层面的重试。Spring Boot 给我们提供了一套内置的重试模板配置起来很简单spring: rabbitmq: listener: simple: retry: enabled: true max-attempts: 3 initial-interval: 1000 multiplier: 2这套配置的含义是消费者处理消息抛出异常时Spring 会进行最多 3 次尝试第一次重试间隔是 1 秒之后每次间隔乘 2也就是 1 秒、2 秒、4 秒。这个机制的好处是重试期间消息不会返回队列不会打扰其他消费者。但这套机制有一个大坑它默认在重试耗尽时还是会走 ack 流程消息直接从队列中丢弃。所以必须配合死信队列或者异常记录表使用确保最终消息不丢失。如果只配置了重试而没有死信队列重试到底后消息就真的没了这在订单类场景是无法接受的。另一种方式是配合手动 nack 和 requeue 参数。在手动确认模式下channel.basicNack(deliveryTag, false, true)把消息重新放回队列配合x-message-ttl可以实现带延时的重试。但这种方式的缺点是消息会打乱顺序而且如果消息本身有毒比如数据格式不对会无限循环重试。所以我在实际项目中更推荐“手动 nack 到死信队列 从死信队列消费做人工补偿”的组合。4.3 死信队列配置与 TTL 参数计算死信队列的概念一条消息在队列中出现下列情况——被消费者拒绝且不重新入队、TTL 过期、队列达到最大长度——就会被转送到指定的死信交换机进而进入死信队列。这相当于消息的“垃圾回收站”但它不是把消息扔掉而是转移到另一个地方让我们继续处理。死信队列最常见的应用是延迟队列。比如订单创建后 30 分钟未支付自动关闭这个场景用死信队列实现就很优雅。实现步骤如下Bean public Queue orderDelayQueue() { return QueueBuilder.durable(order.delay.queue) .withArgument(x-dead-letter-exchange, RabbitConfig.DEAD_EXCHANGE_NAME) .withArgument(x-dead-letter-routing-key, order.close) .withArgument(x-message-ttl, 30 * 60 * 1000) .build(); } Bean public Queue orderCloseQueue() { return QueueBuilder.durable(order.close.queue).build(); }这里的关键是把消息发送到order.delay.queue该队列设置了 TTL 30 分钟消息进入队列后等待 30 分钟过期后被自动转发到死信交换机DEAD_EXCHANGE_NAME路由键为order.close最终进入order.close.queue。消费者监听order.close.queue收到消息时执行关单操作。有个细节必须注意队列级别的 TTL 是消息在哪条队列就统一生效一旦消息过期RabbitMQ 会将其转投到死信交换机。但如果队列中积压了好几条 TTL 不同的消息RabbitMQ 是按“头部消息的到期时间优先”处理的也就是说不保证严格的时间精度只能保证“至少等待 TTL 时长”。如果你需要精确到秒级的延迟建议使用rabbitmq_delayed_message_exchange插件用延迟交换机替代死信队列。5. 常见问题与排查技巧实录5.1 连接超时与端口排查Spring Boot 应用启动后一直报Connection refused这是接触 RabbitMQ 最常遇到的问题。先确认几个环节RabbitMQ 容器或者本机服务是否启动端口是否暴露正确虚拟主机、账号密码是否匹配。之前我在一个团队里遇到过很隐蔽的问题两个环境共用同一个 RabbitMQ但虚拟主机virtual-host不同有人在配置里漏写了 virtual-host导致应用一直连不上。Spring Boot 默认的 virtual-host 是/如果你的账号权限绑定在别的 vhost 下应用当然会拒绝连接。连接端口的问题也需要定位清楚。5672是业务连接端口15672是管理界面端口。有人排查的时候只确认了管理界面能打开就认为 RabbitMQ 运行正常结果应用连不上因为5672端口可能被防火墙拦截或者映射没加。用 telnet 或 nc 命令先测一下两个端口是否都通telnet localhost 5672 telnet localhost 15672端口全通但认证失败时RabbitMQ 的日志会记录具体原因可以在容器里通过docker logs rabbitmq查看。认证失败主要看账号是否存在、密码是否包含特殊字符导致 yml 解析异常、以及用户是否有对应 vhost 的权限。排查这类问题与其在代码里打日志乱猜不如先打开管理界面看看队列和连接的实时状态再回到配置检查账号权限效率高得多。5.2 消息丢失与幂等性保障消息丢失是消息队列场景最致命的问题。RabbitMQ 里消息丢失可能发生在三个环节生产者发送时丢了、Broker 存储时丢了、消费者处理时丢了。生产者的丢失通过开启 Confirm 回调来感知Broker 的丢失通过持久化交换机、队列和消息来解决消费者的丢失则通过手动确认机制来杜绝。消息持久化不是交换机或者队列持久化就完事的。如果发送消息时MessageProperties里的deliveryMode不是PERSISTENT_TEXT_PLAIN消息本身还是非持久化的。在我们convertAndSend方法中如果用的是自定义对象Spring 会在转换后默认设置持久化但如果你手动构造 Message就要显式指定MessageProperties props new MessageProperties(); props.setDeliveryMode(MessageDeliveryMode.PERSISTENT);幂等性则是另一个高频隐患。RabbitMQ 的 at-least-once 投递语义决定了消息有可能被重复消费比如消费者处理完消息后在发送 ack 之前网络闪断Broker 会重新投递这条消息。解决重复消费的思路是在消费端记录消息的唯一标识。我在项目里的做法是生产者在消息头塞入全局唯一 ID消费者收到消息后先查去重表如果 ID 存在说明已处理直接 ack 跳过。RabbitListener(queues RabbitConfig.QUEUE_NAME) public void onMessage(Message message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) { String msgId message.getMessageProperties().getMessageId(); if (duplicateChecker.isDuplicate(msgId)) { channel.basicAck(deliveryTag, false); return; } // ... 执行业务逻辑 duplicateChecker.record(msgId); channel.basicAck(deliveryTag, false); }5.3 消费者线程阻塞与消息积压的排查消息积压是另一个头疼的问题。现象是管理界面的队列里消息数量持续增长消费者不消费或者消费很慢。排查这类问题我一般先从消费者程序入手看日志里是否一直打印异常看线程 dump 里消费者的线程卡在什么地方。遇到过的最典型原因是消费端调用了外部 HTTP 接口而那个接口响应超时导致消费线程全被阻塞在 HTTP 调用上。看起来队列有大量的消息实际消费者的线程都在等响应自然消费不动。解决方案是把请求超时时间调短或者把耗时操作异步化如果用了同步调用一定要给 RestTemplate 或 OpenFeign 设置 connectTimeout 和 readTimeout。还有一些低级的配置错误也会导致不消费。比如消费者监听的队列名称写错了或者程序里的 Bean 没被扫描到RabbitListener 没有生效。遇到这种情况在配置类里打一行日志确认容器工厂是否初始化看看控制台的 RabbitListenerEndpointRegistry 日志就能定位。还有个必须提醒的点在手动确认模式下如果消费者抛出异常后没有 catch 住异常会传播到监听容器Spring 会执行恢复逻辑默认行为是反复重试。如果重试配置不当可能出现“一个异常消息导致消费者线程不断重试后续消息全部无法消费”的雪崩效应。我的处理习惯是消费方法内部 try-catch 全部业务异常业务可重试的抛出异常让重试机制处理业务不可重试的直接记录日志并转入死信队列绝不让异常在监听器外层无限传播。6. 集群部署与生产环境优化经验6.1 Docker 搭建 RabbitMQ 集群要暴露哪些端口生产环境单节点 RabbitMQ 通常不够用至少要做镜像队列集群实现高可用。Docker 搭建集群时端口暴露是关键。每台 RabbitMQ 节点需要暴露的端口有几个层次对外提供服务的5672AMQP和15672管理界面集群节点内部通信的25672集群间通信、4369epmd 节点发现以及如果启用了 MQTT、STOMP 等插件还需要对应端口。搭建时我建议至少保留5672、15672、25672、4369这四个端口的一致性暴露否则集群节点可能互相连不上。集群节点之间的连接在 Docker 网络里要特别注意 hostname 的解析。RabbitMQ 节点名称默认取容器 hostname如果容器重启导致 hostname 变化集群可能失效。简单有效的做法是在启动容器时用-h参数固定 hostname并将节点加入集群时统一用rabbit节点名方式互相连接。这块处理不好费掉的排查时间往往比业务代码开发时间还多。6.2 镜像队列与高可用配置普通集群模式下队列和消息只存在一个节点上其他节点只存储元数据如果那个节点挂了消息就丢失了。解决方法是声明队列时指定镜像策略rabbitmqctl set_policy ha-all ^order\. {ha-mode:all}这条命令把所有以order.开头的队列设置为全节点镜像消息会在集群所有节点上都存一份副本。一个节点挂了其他节点还能继续消费这是目前比较稳妥的高可用方案。读取消息时仍然从队列所在节点读取所以镜像队列在高并发场景下写入会有额外的复制开销因此也不能所有队列都无脑镜像要按业务重要程度分级配置。Spring Boot 消费端连接集群时可以在配置中列出多个地址spring: rabbitmq: addresses: 192.168.1.10:5672,192.168.1.11:5672,192.168.1.12:5672 username: admin password: admin123这样当其中一个节点不可用时客户端会自动尝试连接其他节点对应用来说是透明的。6.3 生产环境的监控与参数调优消息中间件的监控是运维环节的重中之重。我建议至少监控几个核心指标队列积压数量、消费速率、发布确认失败次数、连接数。RabbitMQ 管理界面自带这些指标的图表但线上最好还是接入 Prometheus通过 rabbitmq-exporter 拉取指标再配合告警规则比如“队列积压超过 1 万条持续 5 分钟”就触发告警。调优方面有两个参数值得重点关注。第一个是prefetchCount前面提过要根据单条消息的处理耗时来调整。处理快的消息几十毫秒可以把 prefetch 调大比如 50处理慢的消息几百毫秒以上尽量调小到 5 ~ 10避免消息堆积在消费者内存中。第二个是 RabbitMQ 的vm_memory_high_watermark默认是 0.4也就是内存使用超过 40% 就会触发阻塞发布。如果业务流量大适当调整到 0.6 可以提高吞吐但也要关注容器内存上限否则操作系统 OOM 更麻烦。在实际项目里Spring Boot 集成 RabbitMQ 从简单到复杂会经历几个阶段先能发能收再考虑消息不丢接着做死信重试最后做集群高可用。没有一套配置能适配所有场景关键是理解每个参数背后的原理出了问题能快速定位到是网络、代码、还是配置的问题。我个人的体会是与其一开始追求各种高级特性不如先把手动确认、死信队列、幂等性这老三样做扎实线上出问题的概率会大幅降低。最后再分享一个实用技巧RabbitMQ 的管理界面可以查看每个队列的连接、消费者状态、消息分布情况排查问题的时候先别急着翻代码把管理界面里“Queues 和 Connections”两个标签页过一遍往往问题就清楚了一大半。特别是消息积压的时候看看消费者是处于 idle 状态还是 blocked 状态前者说明没有消息进来或者路由不对后者说明消费者线程卡住了排查方向完全不同。
返回列表