RabbitMQ死信队列延迟消息

wen java案例 1

本文目录导读:

RabbitMQ死信队列延迟消息

  1. 核心概念
  2. 架构图(关键)
  3. 代码示例(Spring Boot + RabbitMQ)
  4. 优点与缺点
  5. 更高级的替代方案

这是一个非常经典且实用的 RabbitMQ 应用场景。死信队列(DLX) + 消息过期时间(TTL) 是 RabbitMQ 实现延迟消息最常见的方式(RabbitMQ 自身没有直接提供类似 delay() 的方法)。

以下是完整的原理、架构组成和代码示例。

核心概念

  1. 延迟消息:消息发送后,不期望消费者立即消费,而是等待一段时间后(30 分钟、1小时)才被投递到消费者。
  2. 死信队列(DLX):当一个消息在队列中变成“死信”后,它可以被重新发送到另一个交换机(死信交换机),最后进入死信队列,被消费者消费。
  3. 消息变为死信的三种情况
    • 消息被消费者拒绝(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 属性,但存在“队列阻塞”问题:如果队列前面的消息没有过期,即使后面消息设置了更长的时间,也不会被消费)。
  • 性能问题:大量消息堆积在普通队列中等待过期,会占用内存。

更高级的替代方案

  1. rabbitmq_delayed_message_exchange 插件(推荐)

    • 原理:使用插件提供的延迟交换机类型 x-delayed-message
    • 优点:支持每个消息独立的延迟时间,性能更好。
    • 引入方式:下载插件,rabbitmq-plugins enable rabbitmq_delayed_message_exchange
  2. Redis Zset 实现

    • 利用有序集合(zset)的 score 存储执行时间戳,配合定时器轮询。
    • 优点:实现简单,性能高。
    • 缺点:Redis 可能丢失消息。
  3. Quartz / Cron Job

    • 定时轮询数据库,找到待处理的任务。
    • 适用:业务数据已经有了,只是需要定时处理。
  • 如果你只需要固定延迟:使用 死信队列 + 队列TTL,配置简单。
  • 如果你需要可变延迟:使用 rabbitmq_delayed_message_exchange 插件Redis Zset
  • 如果对可靠性要求极高:建议使用死信队列方案,并配合消息持久化、手动ACK。

抱歉,评论功能暂时关闭!