5人参与 • 2026-09-21 • Mysql
延迟队列是 rabbitmq 最常被问到的进阶话题。但很多文章只讲一种方案,不说适用场景和坑。这篇把三种方案全部实现一遍,给出完整代码、压测对比和选型建议。建议收藏,做延迟任务时直接对照抄。
假设你在做外卖平台。用户下单后,如果 15 分钟内没有支付,订单要自动取消。
这个需求的核心是:一条消息发出后,不要立刻被消费,而是等 15 分钟再被处理。
这就是延迟队列。
rabbitmq 本身没有"延迟队列"这个队列类型,但可以通过三种方式实现:
每种方案都有适用场景和坑,下面逐一拆解。
这个方案的思路是:消息先进入一个"等待队列",设置 ttl(过期时间),过期后自动转发到死信交换机,再路由到真正的消费队列。

关键点:等待队列不能有消费者。如果有消费者,消息一进队列就被消费了,根本等不到过期。
@configuration
public class delayqueueconfig {
// 等待队列:消息在这里等 15 分钟
@bean
public queue waitqueue() {
return queuebuilder.durable("order.wait.queue")
.ttl(15 * 60 * 1000) // 15 分钟过期
.deadletterexchange("order.dlx.exchange") // 过期后转发到死信交换机
.deadletterroutingkey("order.delay") // 转发时用的路由键
.build();
}
// 死信交换机
@bean
public directexchange dlxexchange() {
return exchangebuilder.directexchange("order.dlx.exchange")
.durable(true)
.build();
}
// 消费队列:真正的消费者在这里
@bean
public queue consumequeue() {
return queuebuilder.durable("order.consume.queue").build();
}
@bean
public binding consumebinding() {
return bindingbuilder.bind(consumequeue())
.to(dlxexchange())
.with("order.delay");
}
}
生产者发送消息到等待队列:
rabbittemplate.convertandsend(
"", // 默认交换机
"order.wait.queue", // 直接发到等待队列
order
);
消费者监听消费队列:
@rabbitlistener(queues = "order.consume.queue")
public void onmessage(order order) {
system.out.println("15 分钟后收到订单: " + order);
// 检查订单是否已支付,未支付则取消
orderservice.cancelifunpaid(order);
}
这个方案有一个致命问题:ttl 只检查队头。
等待队列(fifo):
┌──────────────────────────────────────────────┐
│ 队头:订单 a (ttl=15min) │
│ 订单 b (ttl=5min) ← 不会提前过期! │
│ 订单 c (ttl=30min) │
│ ... │
└──────────────────────────────────────────────┘
如果订单 a 是 15 分钟延迟,订单 b 是 5 分钟延迟,订单 b 必须等订单 a 出队后,才会被检查。也就是说,订单 b 实际延迟了 15 分钟,而不是 5 分钟。
这个问题在业务上往往是不可接受的。
方案一:延迟时间统一。 如果所有消息延迟时间一样,队头阻塞就不存在了。比如所有订单都是 15 分钟后取消。
方案二:按延迟时间分队列。 不同延迟时间的消息走不同的等待队列:
// 5 分钟延迟队列
@bean
public queue waitqueue5min() {
return queuebuilder.durable("order.wait.5min")
.ttl(5 * 60 * 1000)
.deadletterexchange("order.dlx.exchange")
.deadletterroutingkey("order.delay")
.build();
}
// 15 分钟延迟队列
@bean
public queue waitqueue15min() {
return queuebuilder.durable("order.wait.15min")
.ttl(15 * 60 * 1000)
.deadletterexchange("order.dlx.exchange")
.deadletterroutingkey("order.delay")
.build();
}
方案三:不用等待队列,改用延迟插件(见方案二)。
rabbitmq_delayed_message_exchange 是 rabbitmq 官方提供的插件。它引入了一种新的交换机类型 x-delayed-message,消息在交换机层面被延迟,到期后才路由到队列。

与 ttl + dlx 的区别:
没有队头阻塞问题,每条消息的延迟时间互相独立。
# 进入 rabbitmq 容器 docker exec -it rabbitmq bash # 启用插件 rabbitmq-plugins enable rabbitmq_delayed_message_exchange
或者用 docker compose:
services:
rabbitmq:
image: rabbitmq:4-management
environment:
rabbitmq_default_user: guest
rabbitmq_default_pass: guest
volumes:
- ./enabled_plugins:/etc/rabbitmq/enabled_pluginsenabled_plugins 文件内容:
[rabbitmq_management,rabbitmq_delayed_message_exchange].
@configuration
public class delaypluginconfig {
// 声明延迟交换机
@bean
public customexchange delayexchange() {
map<string, object> args = new hashmap<>();
args.put("x-delayed-type", "direct"); // 底层路由类型
return new customexchange(
"order.delay.exchange",
"x-delayed-message", // 延迟交换机类型
true, // durable
false, // autodelete
args
);
}
@bean
public queue delayqueue() {
return queuebuilder.durable("order.delay.queue").build();
}
@bean
public binding delaybinding() {
return bindingbuilder.bind(delayqueue())
.to(delayexchange())
.with("order.delay")
.noargs();
}
}
发送时设置延迟时间:
rabbittemplate.convertandsend(
"order.delay.exchange",
"order.delay",
order,
message -> {
// 设置延迟时间:15 分钟
message.getmessageproperties().setdelaylong(15 * 60 * 1000l);
// 消息持久化
message.getmessageproperties()
.setdeliverymode(messagedeliverymode.persistent);
return message;
}
);
消费者和普通消费者一样:
@rabbitlistener(queues = "order.delay.queue")
public void onmessage(order order) {
system.out.println("15 分钟后收到订单: " + order);
orderservice.cancelifunpaid(order);
}
坑一:不支持非持久化消息。 延迟插件要求消息必须持久化,否则重启后延迟消息会丢。
坑二:延迟时间上限。 延迟时间不能超过 x-max-delay(默认无限制,但实际受限于 erlang timer 的精度和内存)。
坑三:高延迟 + 大吞吐时内存压力大。 延迟消息存在交换机内部,大量消息堆积会占用内存。建议设置合理的延迟时间和消息量。
坑四:集群行为。 在集群中,延迟消息只存在于声明交换机的节点上。如果该节点挂了,延迟消息可能丢失。生产环境建议配合 quorum 队列使用。
rabbitmq streams 是 3.9 引入的新特性,是一种持久化的追加日志,类似 kafka 的 partition。它支持多消费者独立读取、offset 管理和消息回放。

用 streams 实现延迟队列的思路:消费者读取消息时,检查消息的时间戳,如果还没到延迟时间,就等待或跳过,稍后再读。
streams 的延迟实现比较复杂,通常需要配合定时任务或轮询:
// 消费者读取时检查时间戳
@rabbitlistener(queues = "order.stream")
public void onmessage(order order, @header(amqpheaders.received_timestamp) long ts) {
long now = system.currenttimemillis();
long delay = 15 * 60 * 1000l;
if (now - ts < delay) {
// 还没到延迟时间,稍后重试
// 实际实现通常需要把消息重新写回 stream 或等待
return;
}
orderservice.cancelifunpaid(order);
}
对于普通延迟队列需求,streams 过重了。 除非你已经在用 streams,否则不建议为了延迟队列引入它。
| 维度 | ttl + dlx | 延迟插件 | streams |
|---|---|---|---|
| 是否需装插件 | 否 | 是 | 否(3.9+ 内置) |
| 队头阻塞 | 有 | 无 | 无 |
| 支持不同延迟时间 | 不推荐 | 支持 | 支持 |
| 延迟精度 | 秒级 | 毫秒级 | 取决于实现 |
| 吞吐量 | 中 | 中 | 极高 |
| 持久化 | 支持 | 要求持久化 | 支持 |
| 集群高可用 | 好 | 一般 | 好 |
| 适用场景 | 固定延迟 | 灵活延迟 | 超高吞吐 |
| 实现复杂度 | 低 | 低 | 高 |
选型建议
回到外卖场景:
延迟队列是 rabbitmq 最实用的进阶功能之一。三种方案各有适用场景,没有银弹。
选型时先问自己三个问题:
回答完这三个问题,方案自然就出来了。
如果本文对你有帮助,欢迎点赞、收藏、关注。后续可以继续写:
到此这篇关于rabbitmq实现延迟队列的三种方案代码和场景对比的文章就介绍到这了,更多相关rabbitmq延迟队列内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
您想发表意见!!点此发布评论
版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。
发表评论