十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Spring Boot集成RabbitMQ延时队列:死信队列+TTL与延迟插件实战

Spring Boot集成RabbitMQ延时队列:死信队列+TTL与延迟插件实战 做订单系统的时候最让人头疼的就是“下单后30分钟未支付自动关闭”这种需求。一开始我用定时任务轮询订单表SQL倒是好写可真跑到线上就露馅了分钟级扫描有延迟表数据量一大数据库压力直线上升更别提高峰期数据库连接被轮询任务占满的酸爽。后来我把这块全部迁移到 RabbitMQ 延时队列上逻辑清晰了消息也不丢了数据库也轻松了。这篇就把 Spring Boot 结合 RabbitMQ 实现延时队列的完整思路、代码、选型对比和踩坑记录都整理出来适合已经会用 RabbitMQ 基础收发消息、想进一步解决业务延时场景的开发者参考。1. 延时队列到底能解决哪些实际问题先说清楚分布式系统里的“延时”和“定时”到底差在哪。定时任务是纯本地的到点就执行延时队列则把“延时”这个状态变成了一条带时间属性的消息流消息在队列里睡够了才被投递给消费者。业务系统不需要自己维护时间状态机也不会因为服务重启就把该做的事忘掉。实际业务里延时队列最常见的场景就这么几类订单超时未支付自动关闭下单成功发一条延时消息30分钟后进入消费者先查订单状态如果仍是未支付就执行关闭和库存回滚。支付结果补偿查询用户扫码支付后不回调等30秒发一条消息去支付渠道主动拉取支付结果。预约提醒通知用户在系统里预约了明天下午2点的服务提前2小时推送一条提醒短信。活动结束后的自动结算秒杀活动结束后延时1小时把浮动的积分、返现、佣金统一结算入账。缓存过期后的回源刷新缓存 key 快到期之前延时发一条消息预热新数据避免缓存击穿。这类需求如果你坚持用定时任务轮询通常要面对两个麻烦一是轮询周期设短了浪费资源设长了业务体验差二是数据量上来之后每次扫描都要在全表里筛时间字段索引再优化也有性能天花板。延时队列的本质优势在于它把“到点该做什么”交给消息中间件去推动业务侧只要准备好消费者即可。RabbitMQ 实现延时队列主流方案有两类死信队列 TTL消息过期时间方案不需要额外插件依赖 RabbitMQ 自带的消息过期机制和死信转发能力。延时插件方案rabbitmq_delayed_message_exchange官方社区生态提供的延时交换机RabbitMQ 3.8 之后已经支持使用体验更接近“开箱即用”的延时消息。两种方案我都完整落地过各有适用场景后面各用一整章把原理和代码讲透。2. 核心机制拆解消息TTL、死信交换机、死信路由是怎么联动的拿方案一来说如果你不懂 TTL 和 Dead Letter 机制代码写得再漂亮也是空中楼阁因为出了问题你根本不知道往哪个方向排查。先花点时间把这几个基础概念串起来。2.1 消息的两种 TTL 设置方式RabbitMQ 里消息的 TTLTime To Live指的是消息在队列里允许存活的时间超过这个时间还没被消费消息就变成“死信”。TTL 有两种设置粒度队列级 TTL在声明队列时用x-message-ttl参数指定单位是毫秒。队列里的所有消息都遵循同一个过期时间。消息级 TTL发送消息时在消息属性里指定expiration字段单位为毫秒。每条消息可以有自己的过期时间。这两种方式混用的时候RabbitMQ 取的是两者中的较小值。也就是说如果队列设置了 30 分钟消息本身设置了 10 分钟消息实际 10 分钟就过期。2.2 死信交换机DLX与死信路由键DLX Routing Key死信队列的官方概念是当一个队列里的消息变成死信之后它会被重新投递到另一个交换机Dead Letter Exchange再由这个交换机路由到对应的死信队列里。消息变成死信的触发条件常见的有三种消息被消费者手动拒绝basic.reject或basic.nack并且设置了requeuefalse。消息超过设置的 TTL自然过期。队列达到最大长度新消息被投递时最老的消息被丢弃进入死信交换机。我们要用的就是第二种TTL 过期触发死信转发。这里最容易被忽略的是队列和死信交换机之间必须用路由键绑定而且这个路由键和普通业务路由键可以完全不是一个语义。很多人认为“死信路由键必须等于消息原始路由键”这是错的。2.3 一条消息从生产到被延时消费的完整路径为了让你脑子里有完整的链路图我用文字描述一下消息在“死信队列 TTL”方案里走过的每一步生产者发送一条消息通过order.business.exchange交换机路由键order.business.routing进入业务队列order.business.queue。业务队列在声明时指定了x-dead-letter-exchangeorder.dead.exchange、x-dead-letter-routing-keyorder.dead.routing并且设置x-message-ttl180000030分钟。这条消息在业务队列里安安稳稳睡 30 分钟。期间没有任何消费者去动它。30 分钟到了RabbitMQ 把这条消息判定为过期自动把消息重新发布到order.dead.exchange交换机。死信交换机根据路由键order.dead.routing把消息路由到死信队列order.dead.queue。延时消费者监听的正是order.dead.queue此时收到消息触发后续业务逻辑比如关闭订单、释放库存。用大白话说业务队列就是“候车室”消息在里面等到发车时间然后被工作人员引导到真正的检票口死信队列。这个候车室本身不需要消费者它的使命就是扣住消息不让走。2.4 延时精度问题为什么说它是“近似延时”方案一的延时精度有一个天然短板消息过期的检测不是在毫秒粒度实时发生的而是在消息到达队列头部时才会判断是否过期。也就是说如果队列里第一条消息 TTL 是 30 分钟第二条消息 TTL 是 5 分钟那第二条消息不可能在第5分钟被投递它必须等第一条消息出队之后才会被头部扫描实际投递时间可能被拉长到 30 分钟以后。这被称为“队头阻塞”问题在设计方案时一定要提前意识到。如果你只是做“30分钟关单”这种固定延时一个队列只放一种 TTL 消息基本没影响。如果同一个队列里需要承载 10 分钟、30 分钟、1 小时等多种延时等级就需要按延时等级拆分队列或者直接用插件方案。3. Spring Boot 集成基于死信队列 TTL 方案的完整实现这一章是能直接抄作业的硬核部分。我的环境是 Spring Boot 2.7.18 spring-boot-starter-amqp 2.7.18 RabbitMQ 3.12生产环境验证通过。后面的代码你复制到项目里改改业务逻辑就能跑。3.1 环境准备与项目依赖RabbitMQ 的安装部署网上教程一抓一大把我用 Docker 居多测试环境直接跑docker run -d --name rabbitmq \ -p 5672:5672 -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:3.12-management这里有两个细节很多人第一次装会栽进去5672 是 AMQP 协议端口Java 客户端连接用这个15672 是管理控制台端口浏览器访问用这个。默认账号guest只允许 localhost 访问你在服务器上部署之后用 Java 代码远程连接会一直报access_refused所以启动时直接建一个可远程访问的账号最省事。Spring Boot 只需要加一个依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency基础配置写进application.ymlspring: rabbitmq: host: 127.0.0.1 port: 5672 username: admin password: admin123 virtual-host: / publisher-confirm-type: correlated publisher-returns: true listener: simple: acknowledge-mode: manual prefetch: 10publisher-confirm-type: correlated和publisher-returns: true这两个配置是生产环境必须开的后面我会单独讲为什么。3.2 定义交换机、队列、绑定关系的配置类用 Java 配置类把 RabbitMQ 里的各种元素一次性声明好。下面是完整的延迟队列配置业务场景就用“订单下单 30 分钟未支付自动关闭”。package com.example.order.config; import org.springframework.amqp.core.*; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.HashMap; import java.util.Map; Configuration public class OrderDelayQueueConfig { public static final String BUSINESS_EXCHANGE order.business.exchange; public static final String BUSINESS_QUEUE order.business.queue; public static final String BUSINESS_ROUTING_KEY order.business.routing; public static final String DEAD_EXCHANGE order.dead.exchange; public static final String DEAD_QUEUE order.dead.queue; public static final String DEAD_ROUTING_KEY order.dead.routing; /** * 业务交换机负责接收生产者发来的消息 */ Bean public DirectExchange businessExchange() { return new DirectExchange(BUSINESS_EXCHANGE); } /** * 死信交换机负责接收业务队列中过期的消息再转发给死信队列 */ Bean public DirectExchange deadExchange() { return new DirectExchange(DEAD_EXCHANGE); } /** * 业务队列 * 1. 设置 x-message-ttl队列里消息统一存活 30 分钟 * 2. 设置 x-dead-letter-exchange消息过期后转发到死信交换机 * 3. 设置 x-dead-letter-routing-key死信交换机用它路由到死信队列 */ Bean public Queue businessQueue() { MapString, Object args new HashMap(8); // 单条消息过期时间30分钟 args.put(x-message-ttl, 30 * 60 * 1000); // 死信交换机 args.put(x-dead-letter-exchange, DEAD_EXCHANGE); // 死信路由键 args.put(x-dead-letter-routing-key, DEAD_ROUTING_KEY); return new Queue(BUSINESS_QUEUE, true, false, false, args); } /** * 死信队列真正的业务消费者监听这个队列 */ Bean public Queue deadQueue() { return new Queue(DEAD_QUEUE, true); } /** * 业务交换机绑定业务队列 */ Bean public Binding businessBinding() { return BindingBuilder.bind(businessQueue()).to(businessExchange()).with(BUSINESS_ROUTING_KEY); } /** * 死信交换机绑定死信队列 */ Bean public Binding deadBinding() { return BindingBuilder.bind(deadQueue()).to(deadExchange()).with(DEAD_ROUTING_KEY); } }几个值得说的设计点交换机统一用DirectExchange路由键精确匹配好排查问题。FanoutExchange广播模式不适合这种精确转发的场景。new Queue(BUSINESS_QUEUE, true, false, false, args)第一个参数是队列名第二个是持久化第三个是排他第四个是自动删除。排他和自动删除默认都设 false没人用的时候队列也不会悄悄消失。业务队列本身不绑定任何消费者这一点至关重要。如果业务队列绑定了消费者消息一进来就被消费走了TTL 就没有意义。3.3 生产者发送延时消息生产者要做的事非常简单把消息发到业务交换机带上路由键order.business.routing消息自然会在业务队列中等待过期。package com.example.order.producer; import com.example.order.config.OrderDelayQueueConfig; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Component; Component public class OrderMessageSender { private final RabbitTemplate rabbitTemplate; public OrderMessageSender(RabbitTemplate rabbitTemplate) { this.rabbitTemplate rabbitTemplate; } public void sendOrderCloseMessage(Long orderId) { // 这里可以传 JSON 字符串也可以传对象生产环境建议统一序列化为 JSON String messageBody {\orderId\:\ orderId \}; rabbitTemplate.convertAndSend( OrderDelayQueueConfig.BUSINESS_EXCHANGE, OrderDelayQueueConfig.BUSINESS_ROUTING_KEY, messageBody ); } }如果你需要每条消息单独设置不同的延时时间就不能依赖队列级的x-message-ttl而应该在发送时给消息属性设置expiration。public void sendDelayMessage(String messageBody, long delayMillis) { MessagePostProcessor messagePostProcessor message - { message.getMessageProperties().setExpiration(String.valueOf(delayMillis)); // 消息持久化防止 RabbitMQ 重启后消息丢失 message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT); return message; }; rabbitTemplate.convertAndSend( OrderDelayQueueConfig.BUSINESS_EXCHANGE, OrderDelayQueueConfig.BUSINESS_ROUTING_KEY, messageBody, messagePostProcessor ); }注意setExpiration接收的是 String 类型不是 long也别传成 int。我之前看到有人写成setExpiration(1800000)编译直接报错折腾了半天才发现是类型问题。3.4 消费者监听死信队列处理过期消息消费者监听的是order.dead.queue。消息过完 30 分钟被死信交换机转到这个队列之后业务逻辑才在这里真正执行。package com.example.order.consumer; import com.example.order.config.OrderDelayQueueConfig; import com.rabbitmq.client.Channel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; Component public class OrderCloseConsumer { private static final Logger log LoggerFactory.getLogger(OrderCloseConsumer.class); RabbitListener(queues OrderDelayQueueConfig.DEAD_QUEUE) public void handleOrderClose(Message message, Channel channel) throws IOException { long deliveryTag message.getMessageProperties().getDeliveryTag(); try { String body new String(message.getBody(), UTF-8); log.info(收到订单超时关闭消息: {}, body); // 这里写真正的业务逻辑 // 1. 根据 orderId 查询订单状态 // 2. 如果还是未支付执行关闭订单 // 3. 释放库存、记录日志 // 手动确认消息 channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(订单超时关闭处理失败, e); // requeuefalse避免消息无限重新入队 channel.basicNack(deliveryTag, false, false); } } }我在消费者这里用了手动确认模式acknowledge-mode: manual。做延时队列手动确认比自动确认稳妥得多。为什么因为消费端拿到延时消息后要执行的通常是一整套业务逻辑比如关单、调数据库、发通知任何一个环节失败都不该让 RabbitMQ 误以为消息已经处理完。手动确认才能做到“消息处理成功才从队列中移除”。3.5 生产者的发布确认消息不能发出去就不管了延时消息和普通消息最大的不同是消息要在队列里“睡”很久才被处理一旦在沉睡中丢失业务在很长一段时间后才会暴露问题。比如订单关闭消息丢了你可能第二天复盘数据时才发现有一批 30 分钟前就该关闭的订单还挂在未支付状态。所以生产者发送时务必配合 RabbitMQ 的发布确认机制。在application.yml里开了publisher-confirm-type: correlated后可以通过回调确认消息是否被 broker 成功接收。PostConstruct public void init() { rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (ack) { log.info(消息发送成功: {}, correlationData ! null ? correlationData.getId() : null); } else { log.error(消息发送失败: {}, cause); // 这里应该做补偿处理比如重新发送或记录异常表 } }); }网上很多教程不提这个但在生产环境这一步是延时队列的兜底生命线。消息发出前先落一张本地消息表状态标记为“待确认”收到 broker 的 confirm 回调后更新为“已投递”补偿任务定期扫描本地消息表发现超时未确认的消息就重新投递。这套本地消息表 Confirm 的模式能保证消息至少被投递一次。4. 更优雅的官方方案基于延迟插件的实现死信队列 TTL 方案虽然不用装插件但队头阻塞和固定 TTL 的问题确实让它在复杂场景下捉襟见肘。RabbitMQ 官方社区提供了一个延时插件rabbitmq_delayed_message_exchange用起来比死信方案自然得多。4.1 插件安装的两种方式插件本质是一个.ez文件需要放到 RabbitMQ 的 plugins 目录下并启用。版本一定要和 RabbitMQ 主版本对应否则插件加载失败。方式一直接安装到本地 RabbitMQ下载对应版本的插件文件后放到 RabbitMQ 安装目录的plugins目录里然后执行rabbitmq-plugins enable rabbitmq_delayed_message_exchange启用后可以通过rabbitmq-plugins list查看插件状态或者登录管理控制台在 Exchanges 页面新建交换机时看到类型里多了一个x-delayed-message选项说明插件已经生效。方式二Docker 启动时带插件用 Docker 跑 RabbitMQ 的时候官方 management 镜像并不自带这个插件我通常会先起一个容器把插件复制出来再重新构建镜像docker run -d --name rabbitmq-temp rabbitmq:3.12-management docker exec rabbitmq-temp rabbitmq-plugins enable rabbitmq_delayed_message_exchange docker commit rabbitmq-temp rabbitmq:3.12-delayed docker rm -f rabbitmq-temp之后就可以用rabbitmq:3.12-delayed这个镜像启动带插件的 RabbitMQ 了。这种方式对团队协作很友好别人拉取镜像就能用不用各自手动装插件。4.2 声明延时交换机x-delayed-message 类型插件方案的核心是新增了一种交换机类型x-delayed-message。这个交换机接收消息后不会立即路由而是根据消息头里的x-delay参数等待指定时间时间到了再把消息路由到绑定的队列。配置类的写法和普通交换机略有不同它不能直接用DirectExchange、TopicExchange这些 RabbitMQ 内置类型需要声明为CustomExchange。package com.example.order.config; import org.springframework.amqp.core.*; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.HashMap; import java.util.Map; Configuration public class DelayExchangeConfig { public static final String DELAY_EXCHANGE order.delay.exchange; public static final String DELAY_QUEUE order.delay.queue; public static final String DELAY_ROUTING_KEY order.delay.routing; Bean public CustomExchange delayExchange() { MapString, Object args new HashMap(4); // x-delayed-type 表示该延时交换机的内部类型这里用 direct args.put(x-delayed-type, direct); return new CustomExchange(DELAY_EXCHANGE, x-delayed-message, true, false, args); } Bean public Queue delayQueue() { return new Queue(DELAY_QUEUE, true); } Bean public Binding delayBinding() { return BindingBuilder.bind(delayQueue()).to(delayExchange()).with(DELAY_ROUTING_KEY).noargs(); } }x-delayed-type这个参数很关键它的作用是指定延时交换机内部转发消息时采用的路由规则。这里设置成direct消息就会根据路由键精确匹配到队列你也可以设置成topic、fanout完全看你业务需要。4.3 发送带 x-delay 头的延时消息发送消息时通过在消息属性里增加x-delay头来指定延时时间单位是毫秒。这条消息会先被交换机扣下等到时间一到自动转发到绑定的队列。package com.example.order.producer; import com.example.order.config.DelayExchangeConfig; import org.springframework.amqp.AmqpException; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageDeliveryMode; import org.springframework.amqp.core.MessagePostProcessor; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Component; Component public class DelayMessageSender { private final RabbitTemplate rabbitTemplate; public DelayMessageSender(RabbitTemplate rabbitTemplate) { this.rabbitTemplate rabbitTemplate; } public void sendDelayMessage(String messageBody, long delayMillis) { MessagePostProcessor messagePostProcessor new MessagePostProcessor() { Override public Message postProcessMessage(Message message) throws AmqpException { message.getMessageProperties().setHeader(x-delay, delayMillis); message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT); return message; } }; rabbitTemplate.convertAndSend( DelayExchangeConfig.DELAY_EXCHANGE, DelayExchangeConfig.DELAY_ROUTING_KEY, messageBody, messagePostProcessor ); } }注意x-delay的值是 long 类型不是 String这一点和方案一的expiration不一样。4.4 消费者代码与普通队列没有区别插件方案最大的好处是消费者写起来极其简单不需要关心消息是怎么被延时的只需要监听延时交换机绑定的队列即可package com.example.order.consumer; import com.example.order.config.DelayExchangeConfig; import com.rabbitmq.client.Channel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; Component public class DelayOrderConsumer { private static final Logger log LoggerFactory.getLogger(DelayOrderConsumer.class); RabbitListener(queues DelayExchangeConfig.DELAY_QUEUE) public void handleDelayMessage(Message message, Channel channel) throws IOException { long deliveryTag message.getMessageProperties().getDeliveryTag(); try { String body new String(message.getBody(), UTF-8); log.info(收到延时消息: {}, body); // 执行业务逻辑 channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(处理延时消息失败, e); channel.basicNack(deliveryTag, false, false); } } }到这里插件方案的核心链路就完整了。生产端设置x-delay延时交换机扣住消息到点转发给队列消费者正常拉取处理。没有死信队列那一堆配置也没有队头阻塞问题每条消息的延时时间相互独立精度理论上可以达到毫秒级。5. 两种方案怎么选别盲目跟风先看你的实际情况我见过不少人一上来就说“用插件方案简单”结果到了公司环境才发现 RabbitMQ 是运维统一管控的集群根本没权限装插件产品又急着上线最后还得回去老老实实写死信队列。选型不是哪个技术新就选哪个而是看你的部署环境和业务容忍度。从几个维度把两种方案放在一起对比对比维度死信队列 TTL 方案延迟插件方案插件依赖不需要原生支持需要安装 rabbitmq_delayed_message_exchange延时精度受队头阻塞影响只能保证“不小于设定值”每条消息独立计时精度高到点即投递动态延时不灵活队列级 TTL 固定消息级 TTL 可动态但受队头阻塞每条消息可自由设置任意延时时间队列堆积队头阻塞会拖慢后续消息数据保存在交换机中延时到后才进入队列不存在队头阻塞运维成本零额外成本需要维护插件版本集群升级时要同步升级插件适用场景固定延时等级、环境运维管控严格延时时间动态变化、精度要求高这里我把死信队列方案排第一不是没原因的。在生产环境里RabbitMQ 集群通常是基础组件团队统一维护的应用开发人员往往没有权限往 broker 上装插件。如果你们团队能完全掌控 RabbitMQ 部署那插件方案肯定是首选开发效率高逻辑也隐晦性低。如果环境受限死信队列 TTL 也完全能顶住绝大多数业务场景关键是做好延时等级的设计把相同 TTL 的消息归到同一个队列。还有一个实践建议如果最终选了死信队列方案建议不要把x-message-ttl写死在队列参数里而是通过消息级别的expiration设置。这样后面如果产品提出“这笔订单延时 45 分钟关闭”的个性化需求你不需要新增队列只需要在发送消息时指定不同的过期时间即可。不过要注意队头阻塞问题依然存在不同过期时间的消息混在同一个队列短延时的消息会被长延时的消息卡住。所以这种做法只适合对延时精度要求不高的场景。6. 实测过程中的那些坑能救一个是一个说实话网上讲实现原理的教程很多但真正跑生产你就会发现坑全在细节里。这一章我把自己踩过的、帮人排查过的坑都列出来每个都是真实场景不是空谈。6.1 坑一消息持久化没设置RabbitMQ 重启后延时消息全丢这是最隐蔽也最致命的坑。方案一的死信队列逻辑里消息从进入业务队列到被死信消费者处理中间要跨越 30 分钟甚至更久。如果生产者发送消息时没有设置MessageDeliveryMode.PERSISTENT消息只会存在内存里RabbitMQ 服务一重启队列里的所有消息瞬间蒸发业务队列被清空订单永远不会被关闭。Spring Boot 默认情况下convertAndSend发送的消息到底是不是持久化的很多人没验证过。我的建议是在消息投递前显式设置持久化标志message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);同时队列本身声明为持久化队列new Queue(name, true)交换机也声明为持久化。6.2 坑二队头阻塞导致短延时消息被长延时消息卡死这个坑在原理部分提过但我还是要单独拿出来说因为很多人实操后会发现一个诡异的现象我发了 10 分钟后要执行的消息怎么过了 30 分钟才被消费原因就是队头阻塞。队列是 FIFO 结构RabbitMQ 判断消息是否过期的逻辑是只看队头消息的 TTL。如果队头消息 TTL 是 30 分钟排在它后面的消息就算 TTL 只有 5 分钟也得等队头消息出队后RabbitMQ 才会继续判断下一条消息。解决办法有两个不同的延时时间用不同的队列。比如 10 分钟关了队列、30 分钟关单队列、1 小时关单队列各配各的死信交换机互不干扰。直接用插件方案。延时交换机收到消息后是根据每条消息的x-delay单独计时的不按队列排队彻底绕开队头阻塞。6.3 坑三expiration 传参类型错误代码跑不通Spring AMQP 里设置消息过期时间message.getMessageProperties().setExpiration(30000);expiration是 String 类型单位毫秒。但很多人会不自觉地写成数字或者从外部配置读了一个秒值直接塞进来。我见过同事把单位搞混把 10 分钟写成了 10 毫秒结果消息秒变死信消费者瞬间收到一堆本不该现在处理的消息。另外注意expiration的最小取值是 1 毫秒传 0 或负数会导致消息立即过期行为不符合直觉。6.4 坑四死信路由键配置不一致消息在交换机里迷路我排查过一个线上告警业务队列里积压了几万条消息死信队列里一条都没有然后消息不断堆积触发告警。最后查出来是x-dead-letter-routing-key配置成了业务路由键order.business.routing而死信交换机order.dead.exchange上绑定的路由键是order.dead.routing。消息过期后被转投到死信交换机路由键对不上交换机不知道把消息丢给谁最后直接丢弃。排查这类问题有一个高效手段登录 RabbitMQ 管理控制台进入 Exchange 页面点开死信交换机查看 Bindings 列表里实际绑定的路由键。再用测试消息手动投递一次看看消息是否能从死信交换机路由到死信队列。一般几分钟就能定位。6.5 坑五消费端报错之后消息无限 requeue形成死循环手动确认模式下如果消费者处理消息抛了异常并且你在basicNack里设置了requeuetrueRabbitMQ 会把这条消息重新放回原队列然后再次投递给消费者。如果业务逻辑有 bug就会无限循环同一批坏消息反复被消费日志疯狂刷屏CPU 飙升其他正常消息全部排队等着。我的建议是处理类异常比如业务校验不通过使用basicAck确认消息记录日志让消息走完生命周期。系统异常比如数据库连接失败先重试几次仍然失败则basicNack(deliveryTag, false, false)让消息进入死信机制或人工补偿流程不要无限 requeue。6.6 坑六TLS 或 NIO 连接配置不当导致长连接被断开RabbitMQ 客户端默认使用 BIO 连接高并发下默认的 Socket 参数可能会在长时间空闲后导致连接被服务端关闭。如果你在生产者或者消费者里设置了比较长的空闲等待时间比如延时消息场景中消费者可能长时间没有消息可消费一定要检查客户端连接工厂的心跳设置。Bean public ConnectionFactory connectionFactory() { CachingConnectionFactory factory new CachingConnectionFactory(); factory.setHost(127.0.0.1); factory.setPort(5672); factory.setUsername(admin); factory.setPassword(admin123); // 心跳间隔防止长时间空闲被服务端断开建议 30 秒 factory.setRequestedHeartBeat(30); // 发布确认 factory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED); factory.setPublisherReturns(true); return factory; }6.7 坑七消费前校验消息是否真的“到点”了插件方案和死信方案在极端情况下都可能出现消息提前被消费的情况。比如 RabbitMQ 节点发生故障转移、客户端连接被重置、消息被重新投递等。对于订单关闭这种强业务约束的场景消费者拿到消息后不能无脑执行关单应该先把消息里的订单信息拿出来去数据库查一下订单的最新状态如果订单已经是已支付状态说明用户在这 30 分钟内完成了支付直接丢弃消息绝不能把已支付订单关闭。如果订单已经是关闭状态说明消息被重复投递了也直接忽略。只有订单状态仍然为待支付才执行关闭动作。这段逻辑就是“幂等性”的落地。做延时队列如果只想着怎么把消息延时发出去忽略消费端的幂等校验上线后迟早被数据问题坑哭。最后说点个人经验延时队列用到现在我最大的感受是方案选型永远比代码实现重要。如果业务上碰到的是固定延时场景死信队列 TTL 完全够用没必要为了“技术先进”去折腾插件如果延时时间经常动态调整或者消息量大、要求精度高那插件方案能让你少踩队头阻塞的坑。另外无论选哪种方案生产者的 Confirm 机制、消费者的手动 ACK、消费端的幂等校验这三件事一定不能省。消息在队列里睡的那几分钟甚至几十分钟恰恰是整个链路上最脆弱、最容易被忽略的环节。我见过太多系统上线前一切正常、上线后延时消息悄悄丢失的案例等到业务方发现订单没关闭已经是第二天了。所以把监控和日志也一并做好延时消息执行耗时、消费积压数、死信队列长度都拉出来盯上心里才踏实。
返回列表