做后端开发的估计都碰到过这种需求订单超时未支付要自动关闭、用户注册N分钟没激活要发提醒、优惠券到期前要推送……以前的玩法大多是定时任务轮询数据库表扫一圈把到期数据捞出来处理。业务量一大轮询间隔不好定、压力也大还容易有“任务明明到时间了却要等下一轮才被执行”的延迟感。我最近在维护一个老项目技术栈正好卡在Spring Boot 1.4上需要把一批延迟任务从定时轮询迁到消息队列上于是把Spring Boot 1.4连接RabbitMQ、再配合官方延时队列插件这套链路完整撸了一遍。这篇文章就把整个过程写出来包括插件装法、连接配置、交换机声明、消息收发代码以及我踩过的几个坑给准备在Spring Boot 1.x老项目里做延时消息的朋友做个参考。1. 项目背景与整体方案选型1.1 什么时候需要用 RabbitMQ 做延时任务延迟任务在业务里实在太常见了。我这次要处理的场景有三类订单下单后15分钟未支付自动关闭、用户注册后24小时未激活发送提醒短信、优惠券过期前3天推送使用通知。不同业务的延迟时间不同有的是分钟级有的是天级。用定时任务轮询有个很尴尬的问题轮询频率决定了任务的实时性。你10分钟扫一次表那一个本该第2分钟就触发关闭的订单最多要等到第10分钟才被处理你把轮询频率调到1分钟数据库查询压力又上来了特别是订单这种大表每次全表扫描很快就把库拖垮。而且分布式环境下多个实例部署还得考虑任务抢占处理逻辑写起来相当别扭。换成RabbitMQ延时队列之后思路就清晰了业务系统在需要延迟操作的时间点直接把一条“以后处理”的消息丢进队列RabbitMQ负责定时投递给消费者消费者只需要关心消息到了之后干什么。实时性由队列本身保证不存在轮询的“空窗期”服务重启也不丢任务削峰填谷的效果还顺带解决了。至于“订单状态要更新”“短信要发送”这些具体动作全部解耦到消费端生产端和消费端互不干扰。1.2 TTL死信队列 vs 延时插件我为什么选插件RabbitMQ做延迟消息社区里长期有两种主流方案动手前一定要先搞明白它们的适用边界。第一种是TTL死信队列DLX方案。思路是给普通交换机加一个“延时队列”这个队列本身没有消费者消息进入后靠TTL存活指定时间到期后变成死信由死信交换机转发到真正的业务队列业务消费者再从业务队列取消息处理。这个方案不依赖任何插件任何RabbitMQ版本都能用。但代价是配置非常繁琐一种延迟时间就需要单独建一个延时队列万一你的业务有5种延迟时长就得建5个队列加5套绑定关系管理页面上都是密密麻麻的线和交换机。它还有个隐藏问题使用单条消息TTL时如果队列头部是一条长延迟消息后面的短延迟消息会被它堵住必须等头部消息到期才能继续投递这就破坏了延时的精确性。第二种就是官方插件方案也就是这次要用到的rabbitmq_delayed_message_exchange。它的核心思想是引入一种新的交换机类型x-delayed-message消息进入这种交换机后不会立即根据路由键投递而是先把整条消息暂存在内部存储中同时启动定时器等到设置的延迟时间一到交换机再按照参数x-delayed-type指定的路由语义direct、topic、fanout等把消息投给匹配的队列。对客户端来说整个过程完全透明生产端只要在消息头里带上x-delay毫秒值消费端跟消费普通队列消息没有区别。我最终选了插件方案主要基于三点考虑一是配置量小一个延迟交换机加一个队列就能覆盖所有延迟场景新增延迟时间只是改发送时的x-delay参数不用再维护一堆队列二是延迟精确每条消息独立计时短延迟不会被长延迟阻塞三是代码侵入小老项目里接这套东西生产端和消费端改动都不大。代价就是多装一个插件、版本要对齐这个代价在可控范围内。2. 环境准备版本匹配、插件安装与项目依赖2.1 RabbitMQ 版本与延时插件安装Windows / Linux先说最容易被坑的版本问题。rabbitmq_delayed_message_exchange插件不是随便下载一个.ez文件往plugins目录里扔就能用的它跟RabbitMQ主版本有严格的对应关系。我在下载插件前特地去GitHub仓库的Releases页面看了一眼里面每个release都明确标注了支持的RabbitMQ版本区间比如RabbitMQ 3.8.x对应插件3.8.0版3.12.x对应插件3.12.0版。从RabbitMQ 3.13开始这个插件直接集成到官方发行包里了不用再单独下载但依然需要手动执行enable命令启用。我本地开发机是Windows 10安装步骤给读者做个参考。先装Erlang再装RabbitMQ两者版本对应关系建议直接查RabbitMQ官方兼容性表格装反了会出现服务起不来的情况。RabbitMQ装完默认服务已经注册成Windows服务用管理员权限打开命令提示符执行cd C:\Program Files\RabbitMQ Server\rabbitmq_server-3.12.4\sbin rabbitmq-plugins enable rabbitmq_delayed_message_exchange执行完可以再执行rabbitmq-plugins list确认状态看到[e] rabbitmq_delayed_message_exchange说明已启用方括号里的e代表enabled。如果RabbitMQ版本低于3.13需要先去GitHub Releases下载对应版本的.ez文件放到sbin同级目录下的plugins文件夹里再回来执行enable。Linux服务器上也是同理只是路径不同我这次在测试服务器上用的命令是这样的# 以 RabbitMQ 3.12.4 为例 cd /usr/lib/rabbitmq/lib/rabbitmq_server-3.12.4/plugins wget https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases/download/v3.12.0/rabbitmq_delayed_message_exchange-3.12.0.ez rabbitmq-plugins enable rabbitmq_delayed_message_exchange插件启用后还有一步我习惯顺手做掉启用管理端插件。后面排查交换机声明、队列绑定、消息投递情况都要靠它rabbitmq-plugins enable rabbitmq_management装完打开浏览器访问http://127.0.0.1:15672默认账号guest/guest。管理端登录后后面所有验证环节都要用到。这里有个日常开发经常踩的坑如果你之前没用过管理端访问15672发现页面打不开八成是rabbitmq_management插件没有enable。另一个坑是改了RabbitMQ服务端口后Spring连接配置里的端口没同步改导致应用启动时连接不断失败后面排查章节会细说。2.2 Spring Boot 1.4 项目依赖与老项目适配Spring Boot 1.4这个版本年代比较久远它默认的spring-boot-starter-parent父工程指定的Java版本是1.6如果你的开发机只有JDK 1.8必须显式覆盖java.version否则编译会直接失败。我的做法是在pom.xml里的properties块里强制指定同时引入spring-boot-starter-amqp和spring-boot-starter-webparent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version1.4.7.RELEASE/version relativePath/ /parent properties java.version1.8/java.version project.build.sourceEncodingUTF-8/project.build.sourceEncoding /properties dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId scopetest/scope /dependency /dependenciesspring-boot-starter-amqp会传递引入spring-amqp和spring-rabbitSpring Boot 1.4.7.RELEASE对应的Spring AMQP版本是1.6.x这个版本已经支持RabbitListener注解、RabbitAdmin自动声明这些能力做延时消息的收发足够用了。特别提醒一句网上很多写Spring Boot整合RabbitMQ的教程是基于2.x甚至3.x的配置类写法拿过来直接抄到1.4上会报错spring-boot-starter-amqp在2.x里有很多类路径和配置项的变化所以老项目接新东西第一原则是看清楚官方文档里你对应的版本号。依赖引入后启动类加上SpringBootApplication就够了。如果项目里要用RabbitListener注册消费者最好显式加上EnableRabbit注解开启监听注解解析虽然Spring Boot自动配置里也做了这件事但显式声明更稳也方便后来人一眼看懂这个项目启用了RabbitMQ的注解监听机制EnableRabbit SpringBootApplication public class DelayDemoApplication { public static void main(String[] args) { SpringApplication.run(DelayDemoApplication.class, args); } }3. 延时队列核心实现连接配置、交换机声明与消息发送3.1 application.yml 连接参数与可靠性配置Spring Boot 1.4的RabbitMQ自动配置读取spring.rabbitmq前缀下的属性。我这次配置的application.yml如下每一项都是实测可行的注释写在后面spring: application: name: demo-mq-delay rabbitmq: host: 127.0.0.1 port: 5672 username: guest password: guest virtual-host: / requested-heartbeat: 60 connection-timeout: 10000 publisher-confirms: true publisher-returns: true template: mandatory: true listener: simple: acknowledge-mode: manual concurrency: 1 max-concurrency: 5 prefetch: 5逐项说明几个容易被忽视的配置。virtual-host默认是/如果你在RabbitMQ里建了独立虚拟主机这里必须和RabbitMQ里的虚拟主机名保持一致否则连接报404错误。requested-heartbeat是心跳间隔单位秒这个参数对老项目特别重要后面排查章节详细介绍。publisher-confirms和publisher-returns在Spring Boot 1.4里直接配置布尔值就能开启发送方确认和消息回归回调但如果换到Spring Boot 2.x配置就变成了publisher-confirm-type: correlated这个差异升级时要特别注意。listener.simple.acknowledge-mode我配置成manual也就是手动ACK模式意思是消费者拿到消息后业务逻辑处理完必须主动调用channel.basicAck告诉RabbitMQ这条消息我处理完了可以删了。不手动ACK会出现什么结果后面也会讲。prefetch是每个消费者一次性从队列预取的消息数量设成5表示消费者一次最多拿5条处理避免一次性拉太多消息堆积在本地导致处理不过来产生重复。3.2 声明延迟交换机、队列和绑定关系连接配置好了接下来是最核心的部分声明延迟交换机。Spring Boot把RabbitMQ的交换机、队列、绑定关系统一交给RabbitAdmin自动声明你只需要在Spring容器里把它们定义成Bean应用启动时RabbitAdmin会自动到服务端创建这些资源。问题在于延迟交换机不是Spring AMQP里已有的DirectExchange、TopicExchange这些现成类而是一个自定义类型x-delayed-message所以你得用CustomExchange手动创建。延迟交换机在创建时必须传入一个参数x-delayed-type这个参数告诉插件底层按什么路由语义投递消息。我这次用的是direct模式RabbitMQ管理端里看到这个交换机的类型显示为x-delayed-message但实际路由行为跟direct完全一样Configuration public class DelayQueueConfig { public static final String DELAY_EXCHANGE delay.exchange; public static final String DELAY_QUEUE delay.queue; public static final String DELAY_ROUTING_KEY delay.routing; Bean public CustomExchange delayExchange() { MapString, Object args new HashMap(); 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(); } }几个细节值得展开说。CustomExchange的第二个参数是交换机类型字符串必须精确写成x-delayed-message大小写错一点都不行。第三个参数true表示交换机持久化第四个参数false表示消息不是每次都要强制落盘这种配置在服务重启后能自动重建交换机已经存在的延迟消息在RabbitMQ进程恢复后也能重新触发定时投递前提是消息本身以持久化模式发送。DelayQueue的构造方法接收一个布尔值true同样表示队列持久化。绑定这里有点意思对插件延迟交换机做BindingBuilder.bind(...).to(...).with(...)是可行的但实际只有直连交换机、主题交换机、扇形交换机、头交换机这四种原生类型才在Spring AMQP里有对应的绑定支持类。CustomExchange本质上不在这个列表里所以Spring Boot注解贼做法是调用.noargs()完成绑定它会把路由键设置好但不附加额外参数。实测这样绑定到x-delayed-message交换机上是可以正常路由消息的放心用。3.3 消息发送用 x-delay 控制延迟时间生产端发送延迟消息核心就一件事在消息的MessageProperties里塞一个x-delay头值为延迟毫秒数。消息发到delay.exchange交换机后插件看到x-delayed-type是direct就会把消息暂存起来等到延迟时间到了再按路由键delay.routing投递到delay.queue队列。用RabbitTemplate发送时有两种写法我推荐convertAndSend配合MessagePostProcessor这样既能传业务对象进去又能在发送前统一给消息附加延迟参数。代码里我用JSON序列化了业务数据实际项目里你可以换成订单VO、上送短信报文序列化方式栓在Producer内部即可Service public class DelayMessageProducer { Autowired private RabbitTemplate rabbitTemplate; public void sendDelayMessage(String data, long delayMillis) { rabbitTemplate.convertAndSend( DelayQueueConfig.DELAY_EXCHANGE, DelayQueueConfig.DELAY_ROUTING_KEY, data, message - { message.getMessageProperties().setHeader(x-delay, delayMillis); message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT); return message; } ); } }这段代码里的lambda是实现MessagePostProcessor接口在真正把消息发送到RabbitMQ之前拦截到Message对象往它的属性里塞x-delay。delayMillis单位是毫秒比如延迟10秒就传10000L。我在实际项目中写了个验证入口用CommandLineRunner在应用启动后自动发两条延迟不同时间的消息Component public class DelayPublishRunner implements CommandLineRunner { Autowired private DelayMessageProducer producer; Override public void run(String... args) throws Exception { System.out.println(发送时间 LocalDateTime.now()); producer.sendDelayMessage(订单 ORDER-10001 超时关闭, 10000L); producer.sendDelayMessage(预约激活提醒短信, 60000L); } }运行后观察消费者打印的接收时间能明显看到第一条消息10秒后才收到第二条1分钟之后才收到这就验证了整个链路是通的。发送方可靠性这里要补充一个配置我在RabbitTemplate上配置了ConfirmCallback和ReturnCallback这是Spring AMQP 1.6就支持的回调能力。ConfirmCallback在消息被RabbitMQ服务端确认接收后触发ack参数就是确认结果ReturnCallback在消息不可路由时触发比如路由键写错了、交换机名不对会带着replyCode和replyText告诉你原因。实际生产环境建议必配不然消息发丢了你只能干瞪眼Configuration public class RabbitTemplateConfig { Bean public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate new RabbitTemplate(connectionFactory); rabbitTemplate.setMandatory(true); rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (!ack) { System.err.println(消息发送失败cause: cause); } }); rabbitTemplate.setReturnCallback((message, replyCode, replyText, exchange, routingKey) - System.err.println(消息不可达 exchange routingKey replyCode replyText)); return rabbitTemplate; } }有一点要注意setMandatory(true)必须开否则消息不可达时RabbitMQ不会调用ReturnCallback而是静默丢弃消息。4. 消息消费端监听器配置与手动确认4.1 RabbitListener 消费延时消息消费端比生产端简单得多。定义一个组件在方法上打上RabbitListener(queues delay.queue)RabbitMQ监听器容器工厂就会为它创建消费者连接监听指定队列。方法触发时延迟到期的消息正好被插件投递进这个队列你的业务逻辑就在这里执行。我这次用手动ACK模式消费消费者方法带上了Message和Channel两个参数拿到消息后先解析body再执行业务逻辑业务成功后手动basicAck确认。测试时故意让其中一条消息处理失败走basicNack分支看它在队列里的重投情况Component public class DelayMessageConsumer { RabbitListener(queues DelayQueueConfig.DELAY_QUEUE) public void onMessage(Message message, Channel channel) throws Exception { long deliveryTag message.getMessageProperties().getDeliveryTag(); try { String body new String(message.getBody(), UTF-8); System.out.println(接收时间 LocalDateTime.now() deliveryTag deliveryTag 内容 body); // 模拟业务处理失败 if (body.contains(超时关闭)) { throw new RuntimeException(业务处理异常); } channel.basicAck(deliveryTag, false); } catch (Exception e) { System.err.println(业务处理失败消息将重回队列 e.getMessage()); channel.basicNack(deliveryTag, false, true); } } }这里basicAck的第二个参数是multiplefalse表示只确认当前这条消息basicNack的第三个参数是requeuetrue表示消息重新放回队列尾部等待下次消费。注意这里有坑后面排查章节会单独讲。4.2 消费端的并发、重试与 ACK 策略消费端配置看起来就是几行yml但背后涉及的线程模型、重试策略一旦理解不到位生产上会出各种怪问题。Spring AMQP的监听器容器默认是SimpleMessageListenerContainer它内部维护了一个或者多个消费线程。yml里listener.simple.concurrency表示初始并发线程数max-concurrency表示最大并发prefetch表示每个消费者的预取数量。线程数开多了并不能让你的延迟消息更快收到因为这些消息是按时间到了才进入队列消费的早晚取决于RabbitMQ的定时投递而不取决于你有多少消费者在等待。并发数影响的是同时能处理多少条已到期消息以及系统吞吐量上限。手动ACK模式下有个非常关键的机制RabbitMQ会记录未确认的消息数量这个数字受prefetch限制。消费者一次最多预取prefetch条消息如果这些消息迟迟不确认RabbitMQ就不会再给这个消费者分发新消息了。我上面prefetch设的是5意味着消费者这边一旦积压了5条未确认消息后面再多的过期消息也会在队列里干等着。所以手动ACK模式必须保证业务处理速度跟得上消费速度否则消息积压就来了。重试机制也要关注。如果消费端一直抛出异常且basicNack时requeue参数给的是true这条消息会无限循环地进出队列造成消息堆积和消费端日志刷屏。我在老项目里的做法是先判断消息能否重新执行能重新执行的短暂requeue几次次数到了就requeuefalse并配合死信交换机把消息挪走。更精细的做法是用yml里的listener.simple.retry配置让Spring AMQP自己控制重试次数重试耗尽了也不重投而是进入死信队列。生产环境强烈建议把重试次数、死信策略都设计好别让消息无限打回。5. 常见故障排查与实操避坑5.1 插件未启用导致交换机声明失败这是最典型的一个坑。应用启动后RabbitAdmin尝试声明delay.exchange交换机但服务端根本不认识x-delayed-message类型于是你看日志会看到类似这样的报错Channel shutdown: channel error; protocol method: #methodchannel.close(reply-code406, reply-textPRECONDITION_FAILED - inequivalent arg type for exchange delay.exchange in vhost /: received x-delayed-message but current is direct)这个报错有两种原因。第一种就是插件没启用RabbitAdmin没法声明延迟交换机解决办法回第二章去enable插件。第二种情况很容易被忽略RabbitMQ里已经存在了一个同名但类型是direct的交换机比如你之前用原生DirectExchange测试过同一个交换机名现在改成延迟交换机服务端检测到类型不匹配直接拒绝。解决办法是把RabbitMQ管理端里残留的delay.exchange删掉或者把当前配置里的交换机名换一个让它能创建出正确类型的交换机。有个排查小技巧不要只盯着应用日志去RabbitMQ管理端的Exchanges页面看一下。正常情况下列表里delay.exchange的Type那一列显示的是x-delayed-message。如果显示的是direct或者topic说明服务端连插件都没识别到或者同名交换机内容残留了。5.2 clean channel shutdown 与连接不稳定网络上搜索“rabbitmq cause: clean channel shutdown; protocol method”能搜到一堆人遇到这个问题。我在实测中也碰到了场景是测试环境RabbitMQ服务运行一段时间后突然心跳超时Spring端抛出连接断开的异常。这个问题的根因大多数是RabbitMQ服务端认为连接闲置太久主动关闭了channel或connection。解决方向有两个一个是把心跳间隔设置得小一点让客户端更频繁地发送心跳保活我上面yml里写的requested-heartbeat: 60意思就是每60秒发一次心跳如果服务端在超时时间内没有收到任何帧就会判定连接死亡并主动断开另一个是把connection-timeout调大到10000毫秒避免网络抖动时连接建立失败。还有一种场景RabbitMQ装在服务器上应用在本机结果用guest账号远程连接会直接报ACCESS_REFUSED初始化连接时就看到403错误。RabbitMQ从3.x开始默认禁止guest账号非本机访问远程连接必须自己新建用户或者把guest的loopback_users限制去掉。开发环境图省事的话我建议直接建一个专用账号在RabbitMQ管理端Administration里新建用户并赋予权限Spring端yml里的username、password、virtual-host三个配置全部对应改好这样后续权限隔离也好做比改RabbitMQ默认配置安全。5.3 消息丢失、重复投递与持久化经验延时消息这东西最怕的就是“消息没了”和“消息重复”。第一个好解决把发送方确认、交换机持久化、队列持久化、消息持久化整套链路都开起来RabbitMQ在正常情况下不会丢消息。但要注意延迟消息在没有到期投递到队列之前是暂存在插件内部的存储结构里的这个阶段的持久化行为不像队列消息那样透明。我在测试时做过一次模拟发一条延迟1小时的持久化消息然后重启RabbitMQ服务发现重启后消息还能正常投递说明插件的持久化是有效的。不过为了安全起见生产上我还是建议配合发送方确认机制使用至少能知道消息到底是发送成功还是失败。第二个问题更隐蔽。消息重复投递有两个来源一是消费者处理成功但没来得及发送ACKRabbitMQ以为处理失败重新投递一次二是网络瞬时抖动导致连接断开RabbitMQ重新建立连接后把未确认消息再次投递。解决重复问题的核心思路不是让RabbitMQ保证不重复而是让业务处理逻辑具备幂等性。我在订单关闭场景里消费端执行前先去查一遍订单状态已经是CANCEL状态就直接ACK跳过这样即使消息重复了几次也不会对业务数据产生负面影响。5.4 老项目维护的两个额外提醒既然技术栈停在Spring Boot 1.4说明系统多半是历史遗留系统这次接延时队列可能只是其中一个改动。有两件事我想提醒同处境的朋友。第一Spring Boot 1.4和Spring Boot 2.x在RabbitMQ这块配置差异不小最典型的就是发送方确认配置项1.4里是publisher-confirms: trueSpring Boot 2.2及以后变成了publisher-confirm-type: correlated。网上教程注释里写的新版配置直接抄到老版本项目里启动时根本不会生效。建议老项目的维护者手里留存一份自己对应版本的官方配置项对照表别偷懒。第二RabbitMQ服务端如果因为安全漏洞、组件升级等原因要升版本一定要先查Erlang对应关系再查延时插件对应关系。我见过有人把RabbitMQ从3.8升到3.11结果延时插件版本没跟上交换机全部声明失败线上消息整整断了半小时。升级前先在测试环境把插件版本和Erlang版本全部对齐然后再动生产。第三管理端页面能观察到的延迟消息状态比较有限。RabbitMQ管理界面里delay.queue的Ready数在消息未到期时始终是0因为消息还躺在delay.exchange里没投到队列你只看队列页面容易误以为消息丢了。要看延迟消息是否正常入插件重点看Exchanges页面的Delay counters也就是交换机投递计数变化那才是确认消息有没有被插件接收的关键指标。最后分享一个我最实在的体会。延时插件这套方案日常用起来最舒服的就是生产端代码极简一个x-delay头解决所有延迟时间配置问题但前提是安装、版本、声明这三步都踏踏实实走对了。我踩过最大的一次坑就是插件没enableRabbitAdmin声明失败后应用还能启动看上去一切正常结果消息发出去石沉大海最后在管理端确认交换机类型不是x-delayed-message才反应过来。所以无论开发还是上线第一件事永远是在管理端看一眼delay.exchange的类型是不是x-delayed-message。如果是你后面写代码的路就都是顺的。