RabbitMQ 延迟队列实现
一、什么是延迟队列
延迟队列是一种消息队列,消息发送后不会立即被消费,而是在指定的延迟时间后才会投递给消费者。
典型场景:
- 订单超时取消(下单30分钟未支付,自动取消)
- 定时提醒通知
- 失败重试(间隔一定时间后重试)
二、RabbitMQ 实现延迟队列的两种方式
| 方案 | 原理 | 复杂度 | 灵活性 |
|---|---|---|---|
| TTL + 死信队列(DLX) | 消息过期后投递到死信交换机 | 中 | 每个延迟需要独立队列 |
rabbitmq_delayed_message_exchange插件 | 原生延迟交换机 | 低 | 每条消息可设不同延迟 |
三、方案一:TTL + 死信队列(DLX)
3.1 核心概念
┌──────────────────────────────────────────────────────────────┐ │ │ │ Producer ──► 普通交换机 ──► 延迟队列(带TTL, 无消费者) │ │ │ │ │ │ 消息过期 │ │ ▼ │ │ 死信交换机(DLX) │ │ │ │ │ ▼ │ │ 实际消费队列 ←── Consumer │ │ │ └──────────────────────────────────────────────────────────────┘关键点:
- 延迟队列:设置
x-message-ttl,不绑定任何消费者,让消息自然过期 - 死信交换机:延迟队列的
x-dead-letter-exchange,过期消息自动转发到这里 - 实际消费队列:绑定到死信交换机,消费者监听此队列
3.2 关键参数说明
| 参数 | 作用 |
|---|---|
x-message-ttl | 队列级别:队列中所有消息的统一 TTL(毫秒) |
x-expires | 消息级别:单条消息的 TTL(发送时设置expiration属性) |
x-dead-letter-exchange | 指定消息成为死信后投递的目标交换机 |
x-dead-letter-routing-key | 死信投递时使用的 routing key(可选,默认沿用原 routing key) |
3.3 消息成为死信(Dead Letter)的三种情况
- 消息 TTL 过期(最常用)
- 队列达到最大长度(
x-max-length) - 消息被消费者拒绝(
basic.reject或basic.nack且requeue=false)
3.4 Java 代码示例(Spring Boot)
importorg.springframework.amqp.core.*;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;@ConfigurationpublicclassDelayQueueConfig{// 交换机定义publicstaticfinalStringORDER_EXCHANGE="order.exchange";// 业务交换机publicstaticfinalStringDELAY_EXCHANGE="order.delay.exchange";// 死信交换机// 队列定义publicstaticfinalStringDELAY_QUEUE="order.delay.queue";// 延迟队列(无消费者)publicstaticfinalStringDEAD_QUEUE="order.dead.queue";// 实际消费队列// Routing KeypublicstaticfinalStringORDER_ROUTING_KEY="order.create";// 延迟时间:30分钟(毫秒)privatestaticfinalintDELAY_TIME=30*60*1000;// ==================== 业务交换机 ====================@BeanpublicDirectExchangeorderExchange(){returnnewDirectExchange(ORDER_EXCHANGE);}// ==================== 延迟队列:绑定 TTL + 死信交换机 ====================@BeanpublicQueuedelayQueue(){returnQueueBuilder.durable(DELAY_QUEUE).ttl(DELAY_TIME)// 消息存活时间.deadLetterExchange(DELAY_EXCHANGE)// 过期后投递的死信交换机.deadLetterRoutingKey(ORDER_ROUTING_KEY)// 死信投递的 routing key.build();}// 业务交换机 → 延迟队列@BeanpublicBindingdelayBinding(){returnBindingBuilder.bind(delayQueue()).to(orderExchange()).with(ORDER_ROUTING_KEY);}// ==================== 死信交换机 ====================@BeanpublicDirectExchangedelayExchange(){returnnewDirectExchange(DELAY_EXCHANGE);}// ==================== 实际消费队列 ====================@BeanpublicQueuedeadQueue(){returnQueueBuilder.durable(DEAD_QUEUE).build();}// 死信交换机 → 实际消费队列@BeanpublicBindingdeadBinding(){returnBindingBuilder.bind(deadQueue()).to(delayExchange()).with(ORDER_ROUTING_KEY);}}消费者:
importcom.rabbitmq.client.Channel;importorg.springframework.amqp.rabbit.annotation.*;importorg.springframework.stereotype.Component;importjava.io.IOException;@ComponentpublicclassOrderDelayConsumer{@RabbitListener(queues=DelayQueueConfig.DEAD_QUEUE)publicvoidhandleDelayedMessage(Stringmessage,Channelchannel,MessageamqpMessage)throwsIOException{try{System.out.println("收到延迟消息:"+message);// 检查订单是否已支付,未支付则取消channel.basicAck(amqpMessage.getMessageProperties().getDeliveryTag(),false);}catch(Exceptione){channel.basicNack(amqpMessage.getMessageProperties().getDeliveryTag(),false,true);}}}生产者:
@ServicepublicclassOrderService{@AutowiredprivateRabbitTemplaterabbitTemplate;publicvoidcreateOrder(StringorderId){// 1. 保存订单到数据库// ...// 2. 发送延迟消息(30分钟后检查支付状态)rabbitTemplate.convertAndSend(DelayQueueConfig.ORDER_EXCHANGE,DelayQueueConfig.ORDER_ROUTING_KEY,orderId);}}3.5 多个不同延迟时间的处理
如果需要同时支持30分钟取消订单和24小时自动确认收货,需要创建多套队列:
// 30分钟延迟@BeanpublicQueuedelayQueue30Min(){returnQueueBuilder.durable("order.delay.30min.queue").ttl(30*60*1000).deadLetterExchange(DELAY_EXCHANGE).deadLetterRoutingKey("order.cancel").build();}// 24小时延迟@BeanpublicQueuedelayQueue24h(){returnQueueBuilder.durable("order.delay.24h.queue").ttl(24*60*60*1000).deadLetterExchange(DELAY_EXCHANGE).deadLetterRoutingKey("order.confirm").build();}注意:每增加一个延迟级别,就要增加一个队列。如果延迟时间种类很多,建议使用方案二。
四、方案二:Delayed Message Exchange 插件(推荐)
4.1 插件安装
# 1. 下载插件(版本需与 RabbitMQ 匹配)# 下载地址:https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases# 2. 放入 RabbitMQ 插件目录cprabbitmq_delayed_message_exchange-3.12.0.ez\/usr/lib/rabbitmq/plugins/# 3. 启用插件rabbitmq-pluginsenablerabbitmq_delayed_message_exchange4.2 工作流程
┌──────────────────────────────────────────────────────┐ │ │ │ Producer ──► 延迟交换机(x-delayed-message) │ │ │ │ │ │ 根据 header "x-delay" 延迟投递 │ │ ▼ │ │ 实际消费队列 ←── Consumer │ │ │ └──────────────────────────────────────────────────────┘对比方案一,少了一层中转,结构大大简化。
4.3 Java 代码示例(Spring Boot)
importorg.springframework.amqp.core.*;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;importjava.util.HashMap;importjava.util.Map;@ConfigurationpublicclassDelayedMessageConfig{publicstaticfinalStringDELAYED_EXCHANGE="order.delayed.exchange";publicstaticfinalStringDELAYED_QUEUE="order.delayed.queue";publicstaticfinalStringDELAYED_ROUTING_KEY="order.delayed";// ==================== 延迟交换机 ====================@BeanpublicCustomExchangedelayedExchange(){Map<String,Object>args=newHashMap<>();args.put("x-delayed-type","direct");// 底层实际交换类型returnnewCustomExchange(DELAYED_EXCHANGE,"x-delayed-message",// 插件提供的交换机类型true,// 持久化false,// 不自动删除args);}// ==================== 队列 ====================@BeanpublicQueuedelayedQueue(){returnQueueBuilder.durable(DELAYED_QUEUE).build();}@BeanpublicBindingdelayedBinding(){returnBindingBuilder.bind(delayedQueue()).to(delayedExchange()).with(DELAYED_ROUTING_KEY).noargs();}}生产者(每条消息独立设置延迟):
@ServicepublicclassDelayedOrderService{@AutowiredprivateRabbitTemplaterabbitTemplate;publicvoidsendDelayedMessage(StringorderId,intdelayMs){rabbitTemplate.convertAndSend(DelayedMessageConfig.DELAYED_EXCHANGE,DelayedMessageConfig.DELAYED_ROUTING_KEY,orderId,message->{// 通过 header 设置延迟时间(毫秒)message.getMessageProperties().setHeader("x-delay",delayMs);returnmessage;});}}五、两种方案对比
| 对比维度 | TTL + DLX | Delayed Message Plugin |
|---|---|---|
| 安装成本 | 无需插件,原生支持 | 需要安装插件 |
| 架构复杂度 | 高(需要死信交换机中转) | 低(一个交换机搞定) |
| 队列数量 | 每个延迟级别需要单独队列 | 一个队列即可 |
| 延迟粒度 | 队列级别统一(或用消息级 TTL) | 消息级别,每条独立设置 |
| 性能 | TTL 到期时会产生额外投递开销 | 内部使用 Mnesia 表存储,大数据量有瓶颈 |
| 消息顺序 | 同队列 FIFO,先入先过期 | 延迟短的消息可能后发先至 |
| 管理可见性 | 两个队列,一目了然 | 延迟中的消息管理界面不可见 |
| 适用场景 | 延迟级别少且固定的场景 | 延迟时间多样、灵活的场景 |
六、常见问题与注意事项
6.1 消息级 TTL 的"坑"
如果使用消息级别的 TTL(expiration字段),消息在队列中不按过期时间排序,而是按入队顺序。即使队列头部的消息还有10分钟才过期,后面已经过期的消息也不会被投递——必须等头部消息过期或消费后,才会检查下一条。
// ❌ 问题场景:消息A TTL=10分钟,消息B TTL=5秒// 消息A先入队,消息B后入队// 结果:消息B必须等消息A过期后才能被投递// ✅ 解决:不同延迟用不同队列(方案一),或使用插件(方案二)6.2 插件方案的延迟上限
x-delay内部使用int32存储,最大延迟约24.8 天(Integer.MAX_VALUE毫秒 ≈ 24.8天)。
6.3 可靠性保证
- 持久化:队列、交换机、消息都要设置为持久化(
durable=true,delivery_mode=2) - 发送端确认:开启
publisher-confirm确保消息成功到达 - 消费端手动 ACK:处理完业务后再确认,避免消息丢失
# application.ymlspring:rabbitmq:publisher-confirm-type:correlated# 发送端确认publisher-returns:true# 路由失败回调listener:simple:acknowledge-mode:manual# 手动ACK6.4 大量延迟消息的性能考虑
- TTL + DLX 方案:过期的瞬间会有大量消息同时进入死信队列,可能造成瞬时压力。可考虑在 TTL 上加随机偏移量缓解。
- 插件方案:延迟消息存储在 Mnesia 表中,百万级延迟消息时内存开销较大,需做好容量规划。
七、总结
选择建议: ┌─ 延迟级别 ≤ 3 个,且固定不变? ──► TTL + DLX(原生,稳定) │ └─ 延迟级别多 / 每条消息延迟不同? ──► Delayed Message Plugin(灵活,简单)