欢迎光临
我们一直在努力

SpringBoot 第十三篇:RabbitMQ 延迟队列

1. 延迟队列概述

1.1 什么是延迟队列

延迟队列(Delayed Queue)是一种特殊的消息队列,它允许消息在发送后不会立即被消费,而是在指定的延迟时间之后才会被投递给消费者。这种机制在需要定时任务、延迟处理等场景中非常有用。

1.2 延迟队列的应用场景

  • 订单超时取消:用户下单后30分钟内未支付,自动取消订单

  • 延时通知:发送提醒消息,如会议开始前15分钟提醒

  • 重试机制:任务执行失败后,延迟一段时间再重试

  • 定时任务:在指定时间执行特定操作

2. RabbitMQ延迟队列实现方式

2.1 死信队列(DLX)方式

通过消息TTL和死信交换机实现延迟效果。

2.2 插件方式

使用RabbitMQ官方提供的延迟消息插件。

3. 基于死信队列的实现

3.1 环境准备

3.1.1 添加依赖

xml

<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>

3.1.2 配置文件

yaml

spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
virtual-host: /
# 开启消息确认机制
publisher-confirm-type: correlated
publisher-returns: true
listener:
simple:
acknowledge-mode: manual
prefetch: 1

3.2 核心概念理解

3.2.1 TTL(Time To Live)

TTL表示消息的存活时间,有两种设置方式:

  • 队列TTL:队列中所有消息都有相同的过期时间

  • 消息TTL:每条消息可以设置不同的过期时间

3.2.2 死信交换机(DLX)

当消息在队列中变成死信时,它会被重新发送到另一个交换机,这个交换机就是DLX。

消息变成死信的情况:

  • 消息被拒绝且不重新入队

  • 消息过期

  • 队列达到最大长度

3.3 详细代码实现

3.3.1 队列和交换机配置类

java

@Configuration
public class RabbitMQConfig {

// 正常业务交换机
public static final String BUSINESS_EXCHANGE = "business.exchange";
// 正常业务队列
public static final String BUSINESS_QUEUE = "business.queue";
// 正常业务路由键
public static final String BUSINESS_ROUTING_KEY = "business.key";

// 死信交换机
public static final String DLX_EXCHANGE = "dlx.exchange";
// 死信队列
public static final String DLX_QUEUE = "dlx.queue";
// 死信路由键
public static final String DLX_ROUTING_KEY = "dlx.key";

/**
* 声明正常业务交换机
*/
@Bean
public DirectExchange businessExchange() {
return new DirectExchange(BUSINESS_EXCHANGE);
}

/**
* 声明死信交换机
*/
@Bean
public DirectExchange dlxExchange() {
return new DirectExchange(DLX_EXCHANGE);
}

/**
* 声明正常业务队列
* 并绑定死信交换机
*/
@Bean
public Queue businessQueue() {
Map<String, Object> args = new HashMap<>();
// 设置死信交换机
args.put("x-dead-letter-exchange", DLX_EXCHANGE);
// 设置死信路由键
args.put("x-dead-letter-routing-key", DLX_ROUTING_KEY);
// 设置队列消息过期时间(单位:毫秒)
args.put("x-message-ttl", 10000); // 10秒

return QueueBuilder.durable(BUSINESS_QUEUE)
.withArguments(args)
.build();
}

/**
* 声明死信队列
*/
@Bean
public Queue dlxQueue() {
return QueueBuilder.durable(DLX_QUEUE).build();
}

/**
* 绑定正常业务队列到正常业务交换机
*/
@Bean
public Binding businessBinding() {
return BindingBuilder.bind(businessQueue())
.to(businessExchange())
.with(BUSINESS_ROUTING_KEY);
}

/**
* 绑定死信队列到死信交换机
*/
@Bean
public Binding dlxBinding() {
return BindingBuilder.bind(dlxQueue())
.to(dlxExchange())
.with(DLX_ROUTING_KEY);
}
}

3.3.2 消息实体类

java

@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class DelayMessage implements Serializable {

private static final long serialVersionUID = 1L;

/**
* 消息ID
*/
private String messageId;

/**
* 消息内容
*/
private String content;

/**
* 延迟时间(单位:毫秒)
*/
private Long delayTime;

/**
* 创建时间
*/
private LocalDateTime createTime;

/**
* 预期消费时间
*/
private LocalDateTime expectConsumeTime;
}

3.3.3 消息生产者

java

@Component
@Slf4j
public class DelayMessageProducer {

@Autowired
private RabbitTemplate rabbitTemplate;

/**
* 发送延迟消息(使用队列TTL)
*/
public void sendDelayMessage(DelayMessage message) {
try {
// 计算预期消费时间
message.setExpectConsumeTime(
LocalDateTime.now().plusSeconds(message.getDelayTime() / 1000)
);

rabbitTemplate.convertAndSend(
RabbitMQConfig.BUSINESS_EXCHANGE,
RabbitMQConfig.BUSINESS_ROUTING_KEY,
message,
new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
// 设置消息持久化
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
}
}
);

log.info("发送延迟消息成功,消息ID:{},发送时间:{},预期消费时间:{}",
message.getMessageId(),
LocalDateTime.now(),
message.getExpectConsumeTime());

} catch (Exception e) {
log.error("发送延迟消息失败,消息ID:{},错误信息:{}",
message.getMessageId(), e.getMessage());
throw new RuntimeException("发送延迟消息失败", e);
}
}

/**
* 发送自定义TTL的延迟消息(使用消息TTL)
*/
public void sendCustomTTLMessage(DelayMessage message) {
try {
message.setExpectConsumeTime(
LocalDateTime.now().plusSeconds(message.getDelayTime() / 1000)
);

rabbitTemplate.convertAndSend(
RabbitMQConfig.BUSINESS_EXCHANGE,
RabbitMQConfig.BUSINESS_ROUTING_KEY,
message,
new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
// 设置消息级别的TTL
message.getMessageProperties().setExpiration(
String.valueOf(message.getDelayTime())
);
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
}
}
);

log.info("发送自定义TTL消息成功,消息ID:{},延迟时间:{}ms",
message.getMessageId(), message.getDelayTime());

} catch (Exception e) {
log.error("发送自定义TTL消息失败,消息ID:{}", message.getMessageId(), e);
throw new RuntimeException("发送自定义TTL消息失败", e);
}
}
}

3.3.4 消息消费者

java

@Component
@Slf4j
public class DelayMessageConsumer {

/**
* 消费死信队列中的消息(即延迟消息)
*/
@RabbitListener(queues = RabbitMQConfig.DLX_QUEUE)
public void processDelayMessage(DelayMessage message,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) {
try {
log.info("收到延迟消息,消息ID:{},消息内容:{},实际消费时间:{},预期消费时间:{},延迟:{}ms",
message.getMessageId(),
message.getContent(),
LocalDateTime.now(),
message.getExpectConsumeTime(),
Duration.between(message.getExpectConsumeTime(), LocalDateTime.now()).toMillis());

// 模拟业务处理
processBusiness(message);

// 手动确认消息
channel.basicAck(deliveryTag, false);

} catch (Exception e) {
log.error("处理延迟消息失败,消息ID:{}", message.getMessageId(), e);
try {
// 处理失败,拒绝消息并重新入队
channel.basicNack(deliveryTag, false, true);
} catch (IOException ex) {
log.error("拒绝消息失败,消息ID:{}", message.getMessageId(), ex);
}
}
}

/**
* 模拟业务处理
*/
private void processBusiness(DelayMessage message) {
// 这里实现具体的业务逻辑
log.info("处理业务逻辑,消息内容:{}", message.getContent());

// 模拟业务处理时间
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}

3.3.5 控制器类

java

@RestController
@RequestMapping("/delay")
@Slf4j
public class DelayMessageController {

@Autowired
private DelayMessageProducer delayMessageProducer;

/**
* 发送固定延迟时间的消息
*/
@PostMapping("/send-fixed")
public ResponseEntity<String> sendFixedDelayMessage(@RequestParam String content) {
DelayMessage message = DelayMessage.builder()
.messageId(UUID.randomUUID().toString())
.content(content)
.delayTime(10000L) // 固定10秒延迟
.createTime(LocalDateTime.now())
.build();

delayMessageProducer.sendDelayMessage(message);
return ResponseEntity.ok("延迟消息发送成功");
}

/**
* 发送自定义延迟时间的消息
*/
@PostMapping("/send-custom")
public ResponseEntity<String> sendCustomDelayMessage(
@RequestParam String content,
@RequestParam Long delayTime) {

DelayMessage message = DelayMessage.builder()
.messageId(UUID.randomUUID().toString())
.content(content)
.delayTime(delayTime)
.createTime(LocalDateTime.now())
.build();

delayMessageProducer.sendCustomTTLMessage(message);
return ResponseEntity.ok("自定义延迟消息发送成功");
}

/**
* 批量发送延迟消息
*/
@PostMapping("/send-batch")
public ResponseEntity<String> sendBatchDelayMessages() {
for (int i = 1; i <= 5; i++) {
DelayMessage message = DelayMessage.builder()
.messageId(UUID.randomUUID().toString())
.content("批量消息-" + i)
.delayTime(i * 5000L) // 5秒、10秒、15秒…
.createTime(LocalDateTime.now())
.build();

delayMessageProducer.sendCustomTTLMessage(message);
}
return ResponseEntity.ok("批量延迟消息发送成功");
}
}

3.4 高级配置和优化

3.4.1 消息确认配置

java

@Configuration
@Slf4j
public class RabbitMQCallbackConfig {

@Autowired
private RabbitTemplate rabbitTemplate;

@PostConstruct
public void init() {
// 消息发送到Exchange确认回调
rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause) {
if (ack) {
log.info("消息发送到Exchange成功,消息ID:{}",
correlationData != null ? correlationData.getId() : "unknown");
} else {
log.error("消息发送到Exchange失败,消息ID:{},原因:{}",
correlationData != null ? correlationData.getId() : "unknown", cause);
// 这里可以实现重发逻辑
}
}
});

// 消息从Exchange路由到Queue失败回调
rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() {
@Override
public void returnedMessage(ReturnedMessage returned) {
log.error("消息从Exchange路由到Queue失败:交换机:{},路由键:{},响应码:{},响应文本:{}",
returned.getExchange(),
returned.getRoutingKey(),
returned.getReplyCode(),
returned.getReplyText());
// 这里可以实现补偿逻辑
}
});
}
}

3.4.2 多级别延迟队列配置

java

@Configuration
public class MultiLevelDelayConfig {

// 定义多个延迟级别
public static final String DELAY_5S_QUEUE = "delay.5s.queue";
public static final String DELAY_10S_QUEUE = "delay.10s.queue";
public static final String DELAY_30S_QUEUE = "delay.30s.queue";
public static final String DELAY_1M_QUEUE = "delay.1m.queue";

public static final String DELAY_5S_ROUTING_KEY = "delay.5s.key";
public static final String DELAY_10S_ROUTING_KEY = "delay.10s.key";
public static final String DELAY_30S_ROUTING_KEY = "delay.30s.key";
public static final String DELAY_1M_ROUTING_KEY = "delay.1m.key";

/**
* 5秒延迟队列
*/
@Bean
public Queue delay5sQueue() {
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", RabbitMQConfig.DLX_EXCHANGE);
args.put("x-dead-letter-routing-key", RabbitMQConfig.DLX_ROUTING_KEY);
args.put("x-message-ttl", 5000);
return QueueBuilder.durable(DELAY_5S_QUEUE).withArguments(args).build();
}

/**
* 10秒延迟队列
*/
@Bean
public Queue delay10sQueue() {
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", RabbitMQConfig.DLX_EXCHANGE);
args.put("x-dead-letter-routing-key", RabbitMQConfig.DLX_ROUTING_KEY);
args.put("x-message-ttl", 10000);
return QueueBuilder.durable(DELAY_10S_QUEUE).withArguments(args).build();
}

// 其他延迟队列类似…

/**
* 绑定多个延迟队列到业务交换机
*/
@Bean
public Binding binding5s() {
return BindingBuilder.bind(delay5sQueue())
.to(new DirectExchange(RabbitMQConfig.BUSINESS_EXCHANGE))
.with(DELAY_5S_ROUTING_KEY);
}

@Bean
public Binding binding10s() {
return BindingBuilder.bind(delay10sQueue())
.to(new DirectExchange(RabbitMQConfig.BUSINESS_EXCHANGE))
.with(DELAY_10S_ROUTING_KEY);
}

// 其他绑定类似…
}

4. 基于插件的实现方式

4.1 插件安装

4.1.1 下载插件

bash

# 下载延迟消息插件
wget https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases/download/v3.12.0/rabbitmq_delayed_message_exchange-3.12.0.ez

# 将插件复制到RabbitMQ插件目录
cp rabbitmq_delayed_message_exchange-3.12.0.ez /usr/lib/rabbitmq/lib/rabbitmq_server-3.12.0/plugins/

# 启用插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange

# 重启RabbitMQ
systemctl restart rabbitmq-server

4.2 插件方式代码实现

4.2.1 配置类

java

@Configuration
public class DelayedPluginConfig {

public static final String DELAYED_EXCHANGE = "delayed.exchange";
public static final String DELAYED_QUEUE = "delayed.queue";
public static final String DELAYED_ROUTING_KEY = "delayed.key";

/**
* 声明延迟交换机(插件方式)
*/
@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();
}
}

4.2.2 生产者

java

@Component
@Slf4j
public class DelayedPluginProducer {

@Autowired
private RabbitTemplate rabbitTemplate;

/**
* 发送延迟消息(插件方式)
*/
public void sendDelayedMessage(DelayMessage message) {
try {
message.setExpectConsumeTime(
LocalDateTime.now().plusSeconds(message.getDelayTime() / 1000)
);

rabbitTemplate.convertAndSend(
DelayedPluginConfig.DELAYED_EXCHANGE,
DelayedPluginConfig.DELAYED_ROUTING_KEY,
message,
new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
// 设置延迟时间(毫秒)
message.getMessageProperties().setDelay(
Math.toIntExact(message.getDelayTime())
);
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
}
}
);

log.info("发送插件延迟消息成功,消息ID:{},延迟时间:{}ms",
message.getMessageId(), message.getDelayTime());

} catch (Exception e) {
log.error("发送插件延迟消息失败,消息ID:{}", message.getMessageId(), e);
throw new RuntimeException("发送插件延迟消息失败", e);
}
}
}

4.2.3 消费者

java

@Component
@Slf4j
public class DelayedPluginConsumer {

/**
* 消费插件方式的延迟消息
*/
@RabbitListener(queues = DelayedPluginConfig.DELAYED_QUEUE)
public void processDelayedMessage(DelayMessage message,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) {
try {
LocalDateTime now = LocalDateTime.now();
long actualDelay = Duration.between(message.getExpectConsumeTime(), now).toMillis();

log.info("收到插件延迟消息,消息ID:{},预期延迟:{}ms,实际延迟:{}ms,误差:{}ms",
message.getMessageId(),
message.getDelayTime(),
actualDelay,
actualDelay – message.getDelayTime());

// 处理业务
processBusiness(message);

// 确认消息
channel.basicAck(deliveryTag, false);

} catch (Exception e) {
log.error("处理插件延迟消息失败,消息ID:{}", message.getMessageId(), e);
try {
channel.basicNack(deliveryTag, false, true);
} catch (IOException ex) {
log.error("拒绝插件延迟消息失败", ex);
}
}
}

private void processBusiness(DelayMessage message) {
log.info("处理插件延迟消息业务,内容:{}", message.getContent());
// 业务处理逻辑
}
}

5. 两种实现方式的对比

5.1 死信队列方式

优点:

  • 无需安装插件,兼容性好

  • 利用RabbitMQ原生特性

  • 实现相对简单

缺点:

  • 延迟时间不精确(队列头部消息阻塞)

  • 每个延迟级别需要创建单独的队列

  • 资源消耗相对较大

5.2 插件方式

优点:

  • 延迟时间精确

  • 支持任意延迟时间

  • 资源消耗小

缺点:

  • 需要安装插件

  • 插件兼容性问题

  • 生产环境需要额外维护

6. 生产环境最佳实践

6.1 消息持久化

java

@Component
public class MessagePersistenceService {

/**
* 发送持久化消息
*/
public void sendPersistentMessage(RabbitTemplate rabbitTemplate,
String exchange,
String routingKey,
Object message) {
rabbitTemplate.convertAndSend(exchange, routingKey, message, msg -> {
// 设置消息持久化
msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
// 设置消息优先级(如果需要)
msg.getMessageProperties().setPriority(5);
// 设置内容类型
msg.getMessageProperties().setContentType("application/json");
return msg;
});
}
}

6.2 异常处理和重试机制

java

@Configuration
@Slf4j
public class RetryConfig {

/**
* 配置重试机制
*/
@Bean
public MessageRecoverer messageRecoverer() {
return new RejectAndDontRequeueRecoverer();
}

@Bean
public RabbitTemplate.RetryTemplateCustomizer retryTemplateCustomizer() {
return new RabbitTemplate.RetryTemplateCustomizer() {
@Override
public void customize(RetryTemplate retryTemplate) {
retryTemplate.setRetryPolicy(
new SimpleRetryPolicy(3,
Collections.singletonMap(Exception.class, true))
);
retryTemplate.setBackOffPolicy(
new ExponentialBackOffPolicy()
);
}
};
}
}

6.3 监控和告警

java

@Component
@Slf4j
public class RabbitMQMonitor {

@Autowired
private RabbitAdmin rabbitAdmin;

/**
* 监控队列状态
*/
@Scheduled(fixedRate = 60000) // 每分钟执行一次
public void monitorQueues() {
try {
Properties queueProperties = rabbitAdmin.getQueueProperties(
RabbitMQConfig.DLX_QUEUE
);

if (queueProperties != null) {
int messageCount = Integer.parseInt(
queueProperties.get("QUEUE_MESSAGE_COUNT").toString()
);

log.info("死信队列消息数量:{}", messageCount);

// 如果消息堆积严重,发送告警
if (messageCount > 1000) {
sendAlert("死信队列消息堆积严重,当前数量:" + messageCount);
}
}
} catch (Exception e) {
log.error("监控队列状态失败", e);
}
}

private void sendAlert(String message) {
// 实现告警逻辑,如发送邮件、短信等
log.warn("告警:{}", message);
}
}

7. 完整示例和测试

7.1 测试控制器

java

@RestController
@RequestMapping("/test")
@Slf4j
public class TestController {

@Autowired
private DelayMessageProducer delayMessageProducer;

@Autowired
private DelayedPluginProducer delayedPluginProducer;

/**
* 测试死信队列方式
*/
@PostMapping("/dlx")
public ResponseEntity<Map<String, Object>> testDLX(@RequestParam String content,
@RequestParam Long delayTime) {
DelayMessage message = DelayMessage.builder()
.messageId(UUID.randomUUID().toString())
.content(content)
.delayTime(delayTime)
.createTime(LocalDateTime.now())
.build();

delayMessageProducer.sendCustomTTLMessage(message);

Map<String, Object> result = new HashMap<>();
result.put("success", true);
result.put("messageId", message.getMessageId());
result.put("sendTime", LocalDateTime.now());
result.put("expectConsumeTime", message.getExpectConsumeTime());

return ResponseEntity.ok(result);
}

/**
* 测试插件方式
*/
@PostMapping("/plugin")
public ResponseEntity<Map<String, Object>> testPlugin(@RequestParam String content,
@RequestParam Long delayTime) {
DelayMessage message = DelayMessage.builder()
.messageId(UUID.randomUUID().toString())
.content(content)
.delayTime(delayTime)
.createTime(LocalDateTime.now())
.build();

delayedPluginProducer.sendDelayedMessage(message);

Map<String, Object> result = new HashMap<>();
result.put("success", true);
result.put("messageId", message.getMessageId());
result.put("sendTime", LocalDateTime.now());
result.put("expectConsumeTime", message.getExpectConsumeTime());

return ResponseEntity.ok(result);
}

/**
* 性能测试
*/
@PostMapping("/performance")
public ResponseEntity<Map<String, Object>> performanceTest(@RequestParam int count) {
long startTime = System.currentTimeMillis();

for (int i = 0; i < count; i++) {
DelayMessage message = DelayMessage.builder()
.messageId(UUID.randomUUID().toString())
.content("性能测试消息-" + i)
.delayTime(5000L) // 5秒延迟
.createTime(LocalDateTime.now())
.build();

delayedPluginProducer.sendDelayedMessage(message);
}

long endTime = System.currentTimeMillis();

Map<String, Object> result = new HashMap<>();
result.put("success", true);
result.put("totalMessages", count);
result.put("totalTime", endTime – startTime);
result.put("throughput", count * 1000.0 / (endTime – startTime));

return ResponseEntity.ok(result);
}
}

7.2 应用启动类

java

@SpringBootApplication
@EnableScheduling
@EnableRabbit
public class DelayQueueApplication {

public static void main(String[] args) {
SpringApplication.run(DelayQueueApplication.class, args);
}

/**
* 配置JSON序列化
*/
@Bean
public MessageConverter jsonMessageConverter() {
return new Jackson2JsonMessageConverter();
}
}

8. 总结

本文详细介绍了SpringBoot整合RabbitMQ实现延迟队列的两种方式:基于死信队列的方式和基于插件的方式。两种方式各有优缺点,在实际项目中可以根据具体需求选择:

  • 对于延迟时间固定且级别不多的场景,推荐使用死信队列方式

  • 对于延迟时间需要精确控制或延迟级别多的场景,推荐使用插件方式

赞(0)
未经允许不得转载:171主机测评 » SpringBoot 第十三篇:RabbitMQ 延迟队列
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址