it编程 > 数据库 > Mysql

RabbitMQ实现延迟队列的三种方案代码和场景对比

5人参与 2026-09-21 Mysql

延迟队列是 rabbitmq 最常被问到的进阶话题。但很多文章只讲一种方案,不说适用场景和坑。这篇把三种方案全部实现一遍,给出完整代码、压测对比和选型建议。建议收藏,做延迟任务时直接对照抄。

引子:一个真实需求

假设你在做外卖平台。用户下单后,如果 15 分钟内没有支付,订单要自动取消。

这个需求的核心是:一条消息发出后,不要立刻被消费,而是等 15 分钟再被处理。

这就是延迟队列

rabbitmq 本身没有"延迟队列"这个队列类型,但可以通过三种方式实现:

  1. ttl + 死信队列:利用消息过期 + 死信转发;
  2. 延迟交换机插件:安装官方插件,交换机层面延迟;
  3. rabbitmq streams:用 streams 的 offset 和时间戳延迟。

每种方案都有适用场景和坑,下面逐一拆解。

方案一:ttl + 死信队列

原理

这个方案的思路是:消息先进入一个"等待队列",设置 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_plugins

enabled_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

原理

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 最实用的进阶功能之一。三种方案各有适用场景,没有银弹。

选型时先问自己三个问题

  1. 延迟时间固定还是灵活?
  2. 能不能装插件?
  3. 吞吐量要求高不高?

回答完这三个问题,方案自然就出来了。

如果本文对你有帮助,欢迎点赞、收藏、关注。后续可以继续写:

到此这篇关于rabbitmq实现延迟队列的三种方案代码和场景对比的文章就介绍到这了,更多相关rabbitmq延迟队列内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!

(0)

您想发表意见!!点此发布评论

推荐阅读

C++ ORM 多数据库异构混合使用(同一项目同时连接 MySQL、PostgreSQL 和 SQLite)

09-21

MySQL中的函数是什么,有什么作用及应用场景分析

09-21

MySQL写一个简单的存储过程(全过程)

09-21

mysql怎么注册成windows服务详细步骤教程

09-21

MySQL创建简单存储过程的新手入门教程

09-21

Ubuntu Server设置静态ip/固定ip指南

09-20

猜你喜欢

版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。

发表评论