当前位置: 代码网 > it编程>编程语言>正则表达式 > RabbitMQ延迟队列怎么实现?TTL+死信队列和插件方案详解

RabbitMQ延迟队列怎么实现?TTL+死信队列和插件方案详解

2026年08月18日 正则表达式 我要评论
一、什么是延迟队列延迟队列是一种消息队列,消息发送后不会立即被消费,而是在指定的延迟时间后才会投递给消费者。典型场景:订单超时取消(下单30分钟未支付,自动取消)定时提醒通知失败重试(间隔一定时间后重

一、什么是延迟队列

延迟队列是一种消息队列,消息发送后不会立即被消费,而是在指定的延迟时间后才会投递给消费者。

典型场景:

  • 订单超时取消(下单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)的三种情况

  1. 消息 ttl 过期(最常用)
  2. 队列达到最大长度(x-max-length
  3. 消息被消费者拒绝(basic.rejectbasic.nackrequeue=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 + dlxdelayed 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=truedelivery_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(灵活,简单)

以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。

(0)

相关文章:

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

发表评论

验证码:
Copyright © 2017-2026  代码网 保留所有权利. 粤ICP备2024248653号
站长QQ:2386932994 | 联系邮箱:2386932994@qq.com