本文目录导读:

这是一个非常经典且实用的 RabbitMQ 应用场景。死信队列(DLX) + 消息过期时间(TTL) 是 RabbitMQ 实现延迟消息最常见的方式(RabbitMQ 自身没有直接提供类似 delay() 的方法)。
以下是完整的原理、架构组成和代码示例。
核心概念
- 延迟消息:消息发送后,不期望消费者立即消费,而是等待一段时间后(30 分钟、1小时)才被投递到消费者。
- 死信队列(DLX):当一个消息在队列中变成“死信”后,它可以被重新发送到另一个交换机(死信交换机),最后进入死信队列,被消费者消费。
- 消息变为死信的三种情况:
- 消息被消费者拒绝(
basic.reject/basic.nack)且不重新入队。 - 消息过期(设置了 TTL)。
- 队列达到最大长度。
- 消息被消费者拒绝(
利用死信队列实现延迟消息的思路:利用“消息过期”这一特性,人为构造过期,让消息过期后转到死信队列,从而实现“延迟消费”。
架构图(关键)
Producer --> Exchange_A (Normal) --> Queue_A (Normal)
|
| (消息过期后,变成死信)
v
Exchange_B (Dead Letter) --> Queue_B (Dead Letter) --> Consumer
- 核心:
Queue_A不设置消费者,消息发送后,没人消费,只能等着过期。 Queue_A需要设置两个重要参数:x-message-ttl:消息存活时间(毫秒)。x-dead-letter-exchange:当消息过期时,应该发往哪个交换机(死信交换机)。x-dead-letter-routing-key:转发时使用的新路由键(通常是死信队列的路由键)。
代码示例(Spring Boot + RabbitMQ)
配置类(定义交换机和队列)
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMQConfig {
// 1. 普通交换机和队列(用于接收原始消息,并让其过期)
public static final String NORMAL_EXCHANGE = "normal.exchange";
public static final String NORMAL_QUEUE = "normal.queue";
public static final String NORMAL_ROUTING_KEY = "normal.key";
// 2. 死信交换机和队列(用于接收过期消息,真正被消费)
public static final String DEAD_LETTER_EXCHANGE = "dead.letter.exchange";
public static final String DEAD_LETTER_QUEUE = "dead.letter.queue";
public static final String DEAD_LETTER_ROUTING_KEY = "dead.letter.key";
// 创建普通交换机 (Direct类型)
@Bean
public DirectExchange normalExchange() {
return new DirectExchange(NORMAL_EXCHANGE);
}
// 创建死信交换机 (Direct类型)
@Bean
public DirectExchange deadLetterExchange() {
return new DirectExchange(DEAD_LETTER_EXCHANGE);
}
// 创建普通队列 - 关键点:设置死信转发和TTL
@Bean
public Queue normalQueue() {
return QueueBuilder.durable(NORMAL_QUEUE)
// 设置消息过期时间(毫秒),这里假设延迟10秒
.ttl(10000)
// 设置死信交换机
.deadLetterExchange(DEAD_LETTER_EXCHANGE)
// 设置死信路由键(对应死信队列的绑定键)
.deadLetterRoutingKey(DEAD_LETTER_ROUTING_KEY)
.build();
}
// 创建死信队列
@Bean
public Queue deadLetterQueue() {
return QueueBuilder.durable(DEAD_LETTER_QUEUE).build();
}
// 绑定普通交换机和普通队列
@Bean
public Binding normalBinding() {
return BindingBuilder
.bind(normalQueue())
.to(normalExchange())
.with(NORMAL_ROUTING_KEY);
}
// 绑定死信交换机和死信队列
@Bean
public Binding deadLetterBinding() {
return BindingBuilder
.bind(deadLetterQueue())
.to(deadLetterExchange())
.with(DEAD_LETTER_ROUTING_KEY);
}
}
生产者(发送消息)
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class ProducerController {
@Autowired
private RabbitTemplate rabbitTemplate;
@GetMapping("/sendDelay")
public String sendDelayMessage() {
String message = "这是一条延迟10秒的消息:" + System.currentTimeMillis();
rabbitTemplate.convertAndSend(
RabbitMQConfig.NORMAL_EXCHANGE,
RabbitMQConfig.NORMAL_ROUTING_KEY,
message
);
return "消息已发送,10秒后消费";
}
}
消费者(监听死信队列)
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
public class Consumer {
@RabbitListener(queues = RabbitMQConfig.DEAD_LETTER_QUEUE)
public void receiveDeadLetterMessage(String message) {
System.out.println("接收到延迟消息:" + message + ",当前时间:" + System.currentTimeMillis());
}
}
优点与缺点
优点:
- 原生支持:RabbitMQ 本身就支持死信机制,不需要额外插件。
- 可靠性高:消息不会丢失,有持久化机制。
缺点:
- 缺乏灵活性:TTL 是在队列上设置的,所有消息的延迟时间相同,如果需要每条消息的延迟时间不同,实现起来较复杂(需要结合消息的
expiration属性,但存在“队列阻塞”问题:如果队列前面的消息没有过期,即使后面消息设置了更长的时间,也不会被消费)。 - 性能问题:大量消息堆积在普通队列中等待过期,会占用内存。
更高级的替代方案
-
rabbitmq_delayed_message_exchange插件(推荐)- 原理:使用插件提供的延迟交换机类型
x-delayed-message。 - 优点:支持每个消息独立的延迟时间,性能更好。
- 引入方式:下载插件,
rabbitmq-plugins enable rabbitmq_delayed_message_exchange。
- 原理:使用插件提供的延迟交换机类型
-
Redis Zset 实现
- 利用有序集合(
zset)的score存储执行时间戳,配合定时器轮询。 - 优点:实现简单,性能高。
- 缺点:Redis 可能丢失消息。
- 利用有序集合(
-
Quartz / Cron Job
- 定时轮询数据库,找到待处理的任务。
- 适用:业务数据已经有了,只是需要定时处理。
- 如果你只需要固定延迟:使用 死信队列 + 队列TTL,配置简单。
- 如果你需要可变延迟:使用
rabbitmq_delayed_message_exchange插件 或 Redis Zset。 - 如果对可靠性要求极高:建议使用死信队列方案,并配合消息持久化、手动ACK。