RabbitMQ 延迟队列实现

📅 2026/8/11 4:57:31
RabbitMQ 延迟队列实现
RabbitMQ 延迟队列实现一、什么是延迟队列延迟队列是一种消息队列消息发送后不会立即被消费而是在指定的延迟时间后才会投递给消费者。典型场景订单超时取消下单30分钟未支付自动取消定时提醒通知失败重试间隔一定时间后重试二、RabbitMQ 实现延迟队列的两种方式方案原理复杂度灵活性TTL 死信队列DLX消息过期后投递到死信交换机中每个延迟需要独立队列rabbitmq_delayed_message_exchange插件原生延迟交换机低每条消息可设不同延迟三、方案一TTL 死信队列DLX3.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 key3.3 消息成为死信Dead Letter的三种情况消息 TTL 过期最常用队列达到最大长度x-max-length消息被消费者拒绝basic.reject或basic.nack且requeuefalse3.4 Java 代码示例Spring Bootimportorg.springframework.amqp.core.*;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;ConfigurationpublicclassDelayQueueConfig{// 交换机定义publicstaticfinalStringORDER_EXCHANGEorder.exchange;// 业务交换机publicstaticfinalStringDELAY_EXCHANGEorder.delay.exchange;// 死信交换机// 队列定义publicstaticfinalStringDELAY_QUEUEorder.delay.queue;// 延迟队列无消费者publicstaticfinalStringDEAD_QUEUEorder.dead.queue;// 实际消费队列// Routing KeypublicstaticfinalStringORDER_ROUTING_KEYorder.create;// 延迟时间30分钟毫秒privatestaticfinalintDELAY_TIME30*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(queuesDelayQueueConfig.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 Bootimportorg.springframework.amqp.core.*;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;importjava.util.HashMap;importjava.util.Map;ConfigurationpublicclassDelayedMessageConfig{publicstaticfinalStringDELAYED_EXCHANGEorder.delayed.exchange;publicstaticfinalStringDELAYED_QUEUEorder.delayed.queue;publicstaticfinalStringDELAYED_ROUTING_KEYorder.delayed;// 延迟交换机 BeanpublicCustomExchangedelayedExchange(){MapString,ObjectargsnewHashMap();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 DLXDelayed Message Plugin安装成本无需插件原生支持需要安装插件架构复杂度高需要死信交换机中转低一个交换机搞定队列数量每个延迟级别需要单独队列一个队列即可延迟粒度队列级别统一或用消息级 TTL消息级别每条独立设置性能TTL 到期时会产生额外投递开销内部使用 Mnesia 表存储大数据量有瓶颈消息顺序同队列 FIFO先入先过期延迟短的消息可能后发先至管理可见性两个队列一目了然延迟中的消息管理界面不可见适用场景延迟级别少且固定的场景延迟时间多样、灵活的场景六、常见问题与注意事项6.1 消息级 TTL 的坑如果使用消息级别的 TTLexpiration字段消息在队列中不按过期时间排序而是按入队顺序。即使队列头部的消息还有10分钟才过期后面已经过期的消息也不会被投递——必须等头部消息过期或消费后才会检查下一条。// ❌ 问题场景消息A TTL10分钟消息B TTL5秒// 消息A先入队消息B后入队// 结果消息B必须等消息A过期后才能被投递// ✅ 解决不同延迟用不同队列方案一或使用插件方案二6.2 插件方案的延迟上限x-delay内部使用int32存储最大延迟约24.8 天Integer.MAX_VALUE毫秒 ≈ 24.8天。6.3 可靠性保证持久化队列、交换机、消息都要设置为持久化durabletruedelivery_mode2发送端确认开启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灵活简单