一、什么是延迟队列
延迟队列是一种消息队列,消息发送后不会立即被消费,而是在指定的延迟时间后才会投递给消费者。
典型场景:
- 订单超时取消(下单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)
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.bean;
import org.springframework.context.annotation.configuration;
@configuration
public class delayqueueconfig {
// 交换机定义
public static final string order_exchange = "order.exchange"; // 业务交换机
public static final string delay_exchange = "order.delay.exchange"; // 死信交换机
// 队列定义
public static final string delay_queue = "order.delay.queue"; // 延迟队列(无消费者)
public static final string dead_queue = "order.dead.queue"; // 实际消费队列
// routing key
public static final string order_routing_key = "order.create";
// 延迟时间:30分钟(毫秒)
private static final int delay_time = 30 * 60 * 1000;
// ==================== 业务交换机 ====================
@bean
public directexchange orderexchange() {
return new directexchange(order_exchange);
}
// ==================== 延迟队列:绑定 ttl + 死信交换机 ====================
@bean
public queue delayqueue() {
return queuebuilder.durable(delay_queue)
.ttl(delay_time) // 消息存活时间
.deadletterexchange(delay_exchange) // 过期后投递的死信交换机
.deadletterroutingkey(order_routing_key) // 死信投递的 routing key
.build();
}
// 业务交换机 → 延迟队列
@bean
public binding delaybinding() {
return bindingbuilder
.bind(delayqueue())
.to(orderexchange())
.with(order_routing_key);
}
// ==================== 死信交换机 ====================
@bean
public directexchange delayexchange() {
return new directexchange(delay_exchange);
}
// ==================== 实际消费队列 ====================
@bean
public queue deadqueue() {
return queuebuilder.durable(dead_queue).build();
}
// 死信交换机 → 实际消费队列
@bean
public binding deadbinding() {
return bindingbuilder
.bind(deadqueue())
.to(delayexchange())
.with(order_routing_key);
}
}
消费者:
import com.rabbitmq.client.channel;
import org.springframework.amqp.rabbit.annotation.*;
import org.springframework.stereotype.component;
import java.io.ioexception;
@component
public class orderdelayconsumer {
@rabbitlistener(queues = delayqueueconfig.dead_queue)
public void handledelayedmessage(string message, channel channel, message amqpmessage) throws ioexception {
try {
system.out.println("收到延迟消息:" + message);
// 检查订单是否已支付,未支付则取消
channel.basicack(amqpmessage.getmessageproperties().getdeliverytag(), false);
} catch (exception e) {
channel.basicnack(amqpmessage.getmessageproperties().getdeliverytag(), false, true);
}
}
}
生产者:
@service
public class orderservice {
@autowired
private rabbittemplate rabbittemplate;
public void createorder(string orderid) {
// 1. 保存订单到数据库
// ...
// 2. 发送延迟消息(30分钟后检查支付状态)
rabbittemplate.convertandsend(
delayqueueconfig.order_exchange,
delayqueueconfig.order_routing_key,
orderid
);
}
}
3.5 多个不同延迟时间的处理
如果需要同时支持 30分钟取消订单 和 24小时自动确认收货,需要创建多套队列:
// 30分钟延迟
@bean
public queue delayqueue30min() {
return queuebuilder.durable("order.delay.30min.queue")
.ttl(30 * 60 * 1000)
.deadletterexchange(delay_exchange)
.deadletterroutingkey("order.cancel")
.build();
}
// 24小时延迟
@bean
public queue delayqueue24h() {
return queuebuilder.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 插件目录 cp rabbitmq_delayed_message_exchange-3.12.0.ez \ /usr/lib/rabbitmq/plugins/ # 3. 启用插件 rabbitmq-plugins enable rabbitmq_delayed_message_exchange
4.2 工作流程
┌──────────────────────────────────────────────────────┐ │ │ │ producer ──► 延迟交换机(x-delayed-message) │ │ │ │ │ │ 根据 header "x-delay" 延迟投递 │ │ ▼ │ │ 实际消费队列 ←── consumer │ │ │ └──────────────────────────────────────────────────────┘
对比方案一,少了一层中转,结构大大简化。
4.3 java 代码示例(spring boot)
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 delayedmessageconfig {
public static final string delayed_exchange = "order.delayed.exchange";
public static final string delayed_queue = "order.delayed.queue";
public static final string delayed_routing_key = "order.delayed";
// ==================== 延迟交换机 ====================
@bean
public customexchange delayedexchange() {
map<string, object> args = new hashmap<>();
args.put("x-delayed-type", "direct"); // 底层实际交换类型
return new customexchange(
delayed_exchange,
"x-delayed-message", // 插件提供的交换机类型
true, // 持久化
false, // 不自动删除
args
);
}
// ==================== 队列 ====================
@bean
public queue delayedqueue() {
return queuebuilder.durable(delayed_queue).build();
}
@bean
public binding delayedbinding() {
return bindingbuilder
.bind(delayedqueue())
.to(delayedexchange())
.with(delayed_routing_key)
.noargs();
}
}
生产者(每条消息独立设置延迟):
@service
public class delayedorderservice {
@autowired
private rabbittemplate rabbittemplate;
public void senddelayedmessage(string orderid, int delayms) {
rabbittemplate.convertandsend(
delayedmessageconfig.delayed_exchange,
delayedmessageconfig.delayed_routing_key,
orderid,
message -> {
// 通过 header 设置延迟时间(毫秒)
message.getmessageproperties()
.setheader("x-delay", delayms);
return message;
}
);
}
}
五、两种方案对比
| 对比维度 | 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.yml
spring:
rabbitmq:
publisher-confirm-type: correlated # 发送端确认
publisher-returns: true # 路由失败回调
listener:
simple:
acknowledge-mode: manual # 手动ack
6.4 大量延迟消息的性能考虑
- ttl + dlx 方案:过期的瞬间会有大量消息同时进入死信队列,可能造成瞬时压力。可考虑在 ttl 上加随机偏移量缓解。
- 插件方案:延迟消息存储在 mnesia 表中,百万级延迟消息时内存开销较大,需做好容量规划。
七、总结
选择建议: ┌─ 延迟级别 ≤ 3 个,且固定不变? ──► ttl + dlx(原生,稳定) │ └─ 延迟级别多 / 每条消息延迟不同? ──► delayed message plugin(灵活,简单)
以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。
发表评论