java

关注公众号 jb51net

关闭
首页 > 软件编程 > java > RabbitMQ延迟队列实现

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

作者:藤原とラふ店丶

想知道RabbitMQ延迟队列怎么实现吗?本文详细讲解TTL+死信队列和延迟插件两种方案,包括原理、Java代码示例和对比,帮你解决订单超时取消、定时提醒等场景下的延迟消息需求,避开消息级TTL的坑,优化性能

一、什么是延迟队列

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

典型场景:

二、RabbitMQ 实现延迟队列的两种方式

方案原理复杂度灵活性
TTL + 死信队列(DLX)消息过期后投递到死信交换机每个延迟需要独立队列
rabbitmq_delayed_message_exchange 插件原生延迟交换机每条消息可设不同延迟

三、方案一:TTL + 死信队列(DLX)

3.1 核心概念

┌──────────────────────────────────────────────────────────────┐
│                                                              │
│   Producer ──► 普通交换机 ──► 延迟队列(带TTL, 无消费者)      │
│                                    │                         │
│                                    │ 消息过期                 │
│                                    ▼                         │
│                              死信交换机(DLX)                  │
│                                    │                         │
│                                    ▼                         │
│                              实际消费队列 ←── Consumer        │
│                                                              │
└──────────────────────────────────────────────────────────────┘

关键点:

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 可靠性保证

# application.yml
spring:
  rabbitmq:
    publisher-confirm-type: correlated  # 发送端确认
    publisher-returns: true             # 路由失败回调
    listener:
      simple:
        acknowledge-mode: manual        # 手动ACK

6.4 大量延迟消息的性能考虑

七、总结

选择建议:

┌─ 延迟级别 ≤ 3 个,且固定不变? ──► TTL + DLX(原生,稳定)
│
└─ 延迟级别多 / 每条消息延迟不同? ──► Delayed Message Plugin(灵活,简单)

以上为个人经验,希望能给大家一个参考,也希望大家多多支持脚本之家。

您可能感兴趣的文章:
阅读全文