开头

我学习 RabbitMQ 高级特性时,一开始最困惑的问题是:

如果用户下单后 10 分钟不付款,系统怎么自动取消订单?

我最开始以为可以写一个定时任务,每隔几分钟扫一次数据库。但继续学 RabbitMQ 后发现,TTL、死信队列和延迟队列正好能解决这类“消息过一段时间再处理”的问题。

这篇文章主要整理我学习过程中对三个知识点的理解:

  1. TTL:消息为什么会过期
  2. 死信队列:失败消息如何被转移
  3. 延迟队列:如何实现“过一段时间再消费”

一、先介绍背景:为什么需要这些高级特性

在普通消息队列里,生产者发送消息,消费者马上消费:

生产者 -> 交换机 -> 队列 -> 消费者

但实际业务里经常不是这么简单。

比如:

  1. 订单 10 分钟未支付,自动取消
  2. 消费者处理失败,消息不能直接丢
  3. 注册成功 7 天后发送短信
  4. 退款申请超过 24 小时未处理,自动通过

这些问题都和一个核心点有关:

消息不是简单地“立刻消费”,而是要根据时间、失败状态、业务规则进行后续处理。


二、用一个具体案例说明问题

假设我现在有一个订单系统:

用户下单 -> 发送订单消息 -> 等待支付 -> 超时取消订单

我一开始想到的方案是定时任务:

// 每隔一段时间扫描未支付订单
@Scheduled(cron = "0 */5 * * * ?")
public void checkUnpaidOrders() {
    // 查询超过10分钟未支付的订单
    // 修改订单状态为已取消
}

这个方案能用,但问题也很明显:

方案 问题
定时任务扫描数据库 数据量大时压力高
短间隔扫描 数据库频繁被查询
长间隔扫描 订单取消不及时
多服务部署 还要考虑重复执行

所以我开始理解:如果能让消息在 10 分钟后自动进入某个队列,再由消费者处理,就更自然了。


三、分析问题出现的原因

RabbitMQ 默认是“消息到了队列,消费者就可以消费”。

但订单超时这种业务需要的是:

消息先等待一段时间 -> 时间到了 -> 再让消费者处理

这里就涉及三个概念:

概念 解决的问题
TTL 消息或队列里的消息多久过期
死信队列 过期、拒绝、队列满的消息去哪
延迟队列 消息延迟一段时间后再被消费

四、一步步给出解决方案

1. TTL:让消息具有“过期时间”

TTL,全称 Time To Live,表示消息的存活时间。

比如我发送一条消息,希望它 10 秒后过期:

rabbitTemplate.convertAndSend(
    Constant.TTL_EXCHANGE_NAME,
    "",
    "ttl test...",
    message -> {
        message.getMessageProperties().setExpiration("10000");
        return message;
    }
);

这里容易忽略一个点:

setExpiration() 的单位是毫秒,并且参数是字符串。

所以:

setExpiration("10000"); // 10 秒

不是:

setExpiration("10"); // 这其实是 10 毫秒

2. 设置队列 TTL

除了给单条消息设置 TTL,也可以给整个队列设置 TTL:

@Bean("ttlQueue2")
public Queue ttlQueue2() {
    return QueueBuilder
            .durable(Constant.TTL_QUEUE2)
            .ttl(20 * 1000)
            .build();
}

也可以用参数方式:

@Bean("ttlQueue2")
public Queue ttlQueue2() {
    Map<String, Object> arguments = new HashMap<>();
    arguments.put("x-message-ttl", 20000);

    return QueueBuilder
            .durable(Constant.TTL_QUEUE2)
            .withArguments(arguments)
            .build();
}

3. 消息 TTL 和队列 TTL 的区别

对比项 消息 TTL 队列 TTL
设置位置 发送消息时设置 声明队列时设置
粒度 每条消息可以不同 队列内消息统一
参数 expiration x-message-ttl
同时存在时 取较小值 取较小值
常见用途 不同消息不同过期时间 统一过期时间

我一开始以为消息过期后一定会马上删除,后来发现不完全是这样。

尤其是给每条消息单独设置 TTL 时,RabbitMQ 不一定会立刻扫描整个队列删除过期消息,通常会在消息即将投递时再判断。


五、补充代码示例或操作命令

1. TTL 交换机和队列配置

public static final String TTL_QUEUE = "ttl_queue";
public static final String TTL_QUEUE2 = "ttl_queue2";
public static final String TTL_EXCHANGE_NAME = "ttl_exchange";
@Bean("ttlExchange")
public FanoutExchange ttlExchange() {
    return ExchangeBuilder
            .fanoutExchange(Constant.TTL_EXCHANGE_NAME)
            .durable(true)
            .build();
}

@Bean("ttlQueue")
public Queue ttlQueue() {
    return QueueBuilder
            .durable(Constant.TTL_QUEUE)
            .build();
}

@Bean("ttlQueue2")
public Queue ttlQueue2() {
    return QueueBuilder
            .durable(Constant.TTL_QUEUE2)
            .ttl(20 * 1000)
            .build();
}

@Bean("ttlBinding")
public Binding ttlBinding(
        @Qualifier("ttlExchange") FanoutExchange exchange,
        @Qualifier("ttlQueue") Queue queue) {
    return BindingBuilder.bind(queue).to(exchange);
}

@Bean("ttlBinding2")
public Binding ttlBinding2(
        @Qualifier("ttlExchange") FanoutExchange exchange,
        @Qualifier("ttlQueue2") Queue queue) {
    return BindingBuilder.bind(queue).to(exchange);
}

发送消息

@RequestMapping("/ttl1")  
public String ttl1(){  
    System.out.println("ttl1 test...");  
    rabbitTemplate.convertAndSend(Constants.TTL_EXCHANGE, "ttl", "ttl test 10s...", message -> {  
        message.getMessageProperties().setExpiration("10000");//过期时间为10s  
        return message;  
    });  
  
    rabbitTemplate.convertAndSend(Constants.TTL_EXCHANGE, "ttl", "ttl test 30s...", message -> {  
        message.getMessageProperties().setExpiration("30000");//过期时间为30s  
        return message;  
    });  
    return "消息发送成功";  
}

2. 死信队列:让异常消息有地方去

死信消息一般来自三种情况:

  1. 消息过期
  2. 消息被拒绝,并且 requeue=false
  3. 队列达到最大长度

死信队列本质上还是普通队列,只是正常队列配置了死信交换机。

public static final String DLX_EXCHANGE_NAME = "dlx_exchange";
public static final String DLX_QUEUE = "dlx_queue";

public static final String NORMAL_EXCHANGE_NAME = "normal_exchange";
public static final String NORMAL_QUEUE = "normal_queue";
@Bean("normalQueue")
public Queue normalQueue() {
    return QueueBuilder
            .durable(Constant.NORMAL_QUEUE)
            .deadLetterExchange(Constant.DLX_EXCHANGE_NAME)
            .deadLetterRoutingKey("dlx")
            .ttl(10 * 1000)
            .maxLength(10L)
            .build();
}

消费者处理失败时,注意这里要设置 requeue=false

@RabbitListener(queues = Constant.NORMAL_QUEUE)
public void listenNormalQueue(Message message, Channel channel) throws Exception {
    long deliveryTag = message.getMessageProperties().getDeliveryTag();

    try {
        System.out.println("接收到消息:" + new String(message.getBody(), "UTF-8"));

        int result = 3 / 0;

        channel.basicAck(deliveryTag, false);
    } catch (Exception e) {
        channel.basicNack(deliveryTag, false, false);
    }
}

我更推荐 basicNack(deliveryTag, false, false),只拒绝当前消息,避免把前面未确认的消息一起拒绝。


六、插入必要的 Mermaid 图示

图 1:TTL、死信队列、延迟队列之间的关系

用途:先把三个概念的关系串起来,避免把它们当成完全独立的知识点。

正常消费

过期 / 被拒绝 / 队列满

设置 TTL

延迟消费效果

生产者发送消息

普通交换机

普通队列

消息是否正常消费

业务处理完成

成为死信

死信交换机 DLX

死信队列 DLQ

消费者处理死信消息

关键节点解释:

  1. TTL 负责让消息“到时间后过期”
  2. 死信交换机负责接收死信消息
  3. 死信队列负责存储死信消息
  4. 延迟队列可以通过 “TTL + 死信队列” 间接实现

图 2:TTL + 死信队列实现延迟队列

用途:说明为什么消费者不是监听原队列,而是监听死信队列。

消费者 dlx_queue dlx_exchange normal_queue normal_exchange 生产者 消费者 dlx_queue dlx_exchange normal_queue normal_exchange 生产者 消息等待 10 秒 发送消息并设置 TTL=10s 路由到普通队列 TTL 到期后变成死信 根据 routingKey 路由到死信队列 消费死信队列中的消息

关键节点解释:

  1. normal_queue 只是等待队列
  2. 消息过期后进入 dlx_queue
  3. 真正的业务消费者应该监听 dlx_queue

图 3:TTL + DLX 方案的顺序问题

用途:解释为什么先发 20 秒消息,再发 10 秒消息时,10 秒消息可能不会先被消费。

队头20s未过期

20s到期

第1条消息: TTL 20s

normal_queue 队头

第2条消息: TTL 10s

RabbitMQ检查队头

后面的10s消息暂时不会被处理

20s消息进入DLQ

再处理10s消息

关键节点解释:

  1. RabbitMQ 会优先看队头消息
  2. 如果队头消息延迟时间更长,后面的短延迟消息可能被挡住
  3. 所以 TTL + DLX 更适合固定延迟场景

七、说明验证方式

1. 验证 TTL

调用接口:

http://127.0.0.1:8080/producer/ttl1

观察 RabbitMQ 管理页面:

  1. 刚发送时,队列 Ready 数量变为 2
  2. 到达 TTL 时间后,消息消失
  3. 如果队列设置了 TTL,队列 Features 中会出现 TTL

这里我分别给普通队列和一个TTL队列发送两条存活时间分别为10s 和 20s 的两条消息,在控制台观察消息的存活状态
image.png

2. 验证死信队列

调用接口:

http://127.0.0.1:8080/producer/dl

发送ttl消息给normal队列

@RequestMapping("/dl")  
public String dl(){  
    System.out.println("dl...");  
    rabbitTemplate.convertAndSend(Constants.NORMAL_EXCHANGE,"normal", "dl test...",message -> {  
        message.getMessageProperties().setExpiration("10000");  
        return message;  
    });  
    System.out.printf("%s 消息发送成功 \n",new Date());  
    return "消息发送成功";  
}

那么 10 秒后,消息会从普通队列进入死信队列。

image.png

通过这种方式,就可以给订单设置TTL,真正的发货系统只需要监听死信队列即可.

管理页面中可以关注:

标识 含义
D durable,持久化
TTL 队列设置了过期时间
Lim 队列设置了最大长度
DLX 设置了死信交换机
DLK 设置了死信 RoutingKey

3. 验证延迟队列插件

通过TTL + 死信的方式可以模拟实现延迟队列的功能, 虽然我们无法直接创建一个延迟队列,但是通过安装插件可以解决这一点
如果使用 RabbitMQ 延迟插件,可以先安装插件:

image.png

把自己的版本信息告诉ai,让ai指导你完成安装即可, 这里不过多介绍

插件安装完成后, 再次打开Web控制台, 可以观察到选项新增x-delayed-message

image.png

Spring AMQP 配置延迟交换机:

import com.amadeus.rabbitextensiondemo.constant.Constants;  
import org.springframework.amqp.core.*;  
import org.springframework.beans.factory.annotation.Qualifier;  
import org.springframework.context.annotation.Bean;  
  
public class DelayConfig {  
    @Bean("dealayQueue")  
    public Queue dealayQueue(){  
        return QueueBuilder  
                .durable(Constants.DELAY_QUEUE)  
                .build();  
    }  
  
    @Bean("delayExchange")  
    public Exchange delayExchange(){  
        return ExchangeBuilder  
                .directExchange(Constants.DELAY_EXCHANGE)  
                .durable(true)  
                .delayed()  
                .build();  
    }  
      
    @Bean("delayBlinding")  
    public Binding delayBinding(@Qualifier("delayQueue") Queue queue, @Qualifier("delayExchange") Exchange exchange){  
        return BindingBuilder  
                .bind(queue)  
                .to(exchange)  
                .with("delay")  
                .noargs();  
    }  
}

发送延迟消息:

rabbitTemplate.convertAndSend(
        Constant.DELAYED_EXCHANGE_NAME,
        "delayed",
        "delayed test 10s",
        message -> {
            message.getMessageProperties().setDelayLong(10000L);
            return message;
        }
);

配置消费者

@Component  
public class DelayListener {  
    @RabbitListener(queues = Constants.DELAY_QUEUE)  
    public void delayHandMessage(Message message, Channel channel) throws Exception {  
        //消费者逻辑  
        System.out.printf("[delay.queue] %tc 接收到消息: %s \n", new Date(), new String(message.getBody(),"UTF-8"));  
    }  
}

image.png

由于我这里没有给延迟队列手动ack,所以在控制台中可以看到未unacked消息记录,当然也可以手动ack

@Component  
public class DelayListener {  
    @RabbitListener(queues = Constants.DELAY_QUEUE)  
    public void delayHandMessage(Message message, Channel channel) throws Exception {  
          
        long deliverTag = message.getMessageProperties().getDeliveryTag();  
        try{  
            //消费者逻辑  
            System.out.printf("[delay.queue] %tc 接收到消息: %s \n",  
                    new Date(),   
new String(message.getBody(),"UTF-8"));  
        }catch (Exception e){  
            System.out.println("消费失败");  
            channel.basicNack(deliverTag,false,true);  
        }  
    }  
}

延迟队列插件方案的好处是:

即使先发送 20 秒消息,再发送 10 秒消息,10 秒消息也可以先到达消费者。不用关心在TTL + 死信的方案下, 由于TTL 长的消息先发导致TTL 短的消息不能即使过期处理导致的问题


八、总结容易踩坑的地方

  1. TTL 单位是毫秒,不是秒。
  2. setExpiration() 传的是字符串,比如 "10000"
  3. 消息 TTL 和队列 TTL 同时存在时,取较小值。
  4. 死信队列不是特殊队列,本质还是普通队列。
  5. 要进入死信队列,普通队列必须配置 x-dead-letter-exchange
  6. 消费失败后如果想进死信队列,requeue 必须是 false
  7. TTL + 死信队列实现延迟队列时,不适合大量不同延迟时间的消息混在一个队列。
  8. 延迟插件需要版本匹配,否则可能启用失败。
  9. RabbitMQ 中已经存在的队列参数不能随便改,参数变化后可能需要删除旧队列或换队列名。

九、几套可以直接复用的提示词/模板

模板 1:TTL 队列配置模板

@Bean("ttlQueue")
public Queue ttlQueue() {
    return QueueBuilder
            .durable("ttl_queue")
            .ttl(10 * 1000)
            .build();
}

模板 2:死信队列配置模板

@Bean("normalQueue")
public Queue normalQueue() {
    return QueueBuilder
            .durable("normal_queue")
            .deadLetterExchange("dlx_exchange")
            .deadLetterRoutingKey("dlx")
            .ttl(10 * 1000)
            .maxLength(10L)
            .build();
}

模板 3:延迟插件发送消息模板

rabbitTemplate.convertAndSend(
        "delayed_exchange",
        "delayed",
        "message body",
        message -> {
            message.getMessageProperties().setDelayLong(10000L);
            return message;
        }
);

模板 4:RabbitMQ 问题排查提示词

我正在排查 RabbitMQ 延迟队列问题。
现象是:
预期是:
当前交换机类型是:
队列参数包括:
消息 TTL 是:
消费者监听的队列是:
请帮我判断消息为什么没有按预期进入消费者。

模板 5:面试回答模板

RabbitMQ 本身普通队列不直接支持延迟消费。
常见实现方式有两种:

第一种是 TTL + 死信队列。
消息先进入普通队列,等待 TTL 到期后变成死信,再通过死信交换机进入死信队列,消费者监听死信队列完成延迟消费。

第二种是使用 rabbitmq_delayed_message_exchange 插件。
它可以直接给消息设置 delay 时间,避免 TTL + DLX 方案中队头阻塞导致的顺序问题。

如果业务延迟时间固定,可以用 TTL + DLX。
如果每条消息延迟时间不同,更推荐使用延迟插件。

结尾总结

这次学完 TTL、死信队列和延迟队列后,我对 RabbitMQ 的理解更清楚了:

  1. TTL 解决的是“消息多久过期”的问题。
  2. 死信队列解决的是“异常消息去哪处理”的问题。
  3. 延迟队列解决的是“消息过一段时间再消费”的问题。
  4. TTL + 死信队列可以实现延迟效果,但有队头阻塞问题。
  5. 延迟插件更适合每条消息延迟时间不同的场景。
  6. 消费失败要进入死信队列,关键是 requeue=false
  7. 学这些高级特性时,不要只记概念,最好用订单超时、退款超时这种业务场景去理解。
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐