当前位置:

首页 > 编程开发 > Java代码实现RabbitMQ延迟消息队列的方法

Java代码实现RabbitMQ延迟消息队列的方法

RabbitMQ延时队列介绍RabbitMQ延时队列是指消息在发送到队列后,并不立即被消费者消费,而是等待一段时间后再被消费者消费。这种队列通常用于实现定时任务,例如,订单超时未支付系统取消订单释放所占库存等。RabbitMQ实现延时队列的方法有多种,其中比较常见的是使用插件或者通过DLX(DeadLetterExchange)机制实现。使用插件实现延时队列RabbitMQ提供了rabbitmq_delayed_message_exchange插件,可以通过该插件实现延时队列。该插件的原理是在消息发送时,

    RabbitMQ 延时队列介绍

    RabbitMQ 延时队列是指消息在发送到队列后,并不立即被消费者消费,而是等待一段时间后再被消费者消费。这种队列通常用于实现定时任务,例如,订单超时未支付系统取消订单释放所占库存等。

    RabbitMQ实现延时队列的方法有多种,其中比较常见的是使用插件或者通过DLX(Dead Letter Exchange)机制实现。

    使用插件实现延时队列

    RabbitMQ提供了rabbitmq_delayed_message_exchange插件,可以通过该插件实现延时队列。该插件的原理是在消息发送时,将消息发送到一个特定的Exchange中,然后该Exchange会根据消息中的延时时间将消息转发到指定的队列中,从而实现延时队列的功能

    使用该插件需要先安装插件,然后创建一个Exchange,并将该Exchange的类型设置为x-delayed-message,然后将该Exchange与队列绑定即可。

    使用DLX机制实现延时队列

    消息的TTL就是消息的存活时间。RabbitMQ可以对队列和消息分别设置TTL。而对队列设置就是队列没有消费者连着的保留时间,也可以对每一个单独的消息做单独的 设置。超过了这个时间,我们认为这个消息就死了,称之为死信。如果队列设置了,消息也设置了,那么会取小的。所以一个消息如果被路由到不同的队 列中,这个消息死亡的时间有可能不一样(不同的队列设置)。这里单讲单个消息的TTL,因为它才是实现延迟任务的关键。可以通过设置消息的expiration字段或者x- message-ttl属性来设置时间,两者是一样的效果

    DLX机制是RabbitMQ提供的一种消息转发机制,它可以将无法被处理的消息转发到指定的Exchange中,从而实现消息的延时处理。具体实现步骤如下:

    • 创建一个普通的Exchange和Queue,并将它们绑定在一起。

    • 创建一个DLX Exchange,并将普通Exchange绑定到该DLX Exchange上。

    • 将Queue设置为具有TTL(Time To Live)属性,并设置消息过期时间。

    • 将Queue绑定到DLX Exchange上。

    当消息过期后,会被发送到DLX Exchange中,然后再由DLX Exchange将消息转发到指定的Exchange中,从而实现延时队列的功能。

    使用DLX机制实现延时队列的优点是不需要安装额外的插件,但是需要对消息的过期时间进行精确控制,否则可能会出现消息过期时间不准确的情况。

    Java语言设置延时队列

    下面是使用 Java 语言通过 RabbitMQ 设置延时队列的步骤:

    安装插件

    首先,需要安装 rabbitmq_delayed_message_exchange 插件。可以通过以下命令安装:

    rabbitmq-plugins enable rabbitmq_delayed_message_exchange

    创建延时交换机

    延时队列需要使用延时交换机。可以使用 x-delayed-message 类型创建一个延时交换机。以下是创建延时交换机的示例代码:

    Map args = new HashMap<>();
    args.put("x-delayed-type", "direct");
    channel.exchangeDeclare("delayed-exchange", "x-delayed-message", true, false, args);

    创建延时队列

    创建延时队列时,需要将队列绑定到延时交换机上,并设置队列的 TTL(Time To Live)参数。以下是创建延时队列的示例代码:

    Map args = new HashMap<>();
    args.put("x-dead-letter-exchange", "delayed-exchange");
    args.put("x-dead-letter-routing-key", "delayed-queue");
    args.put("x-message-ttl", 5000);
    channel.queueDeclare("delayed-queue", true, false, false, args);
    channel.queueBind("delayed-queue", "delayed-exchange", "delayed-queue");

    在上述代码中,将队列绑定到延时交换机上,并设置了队列的 TTL 参数为 5000 毫秒,即消息在发送到队列后,如果在 5000 毫秒内没有被消费者消费,则会被转发到 delayed-exchange 交换机上,并发送到 delayed-queue 队列中。

    发送延时消息

    发送延时消息时,需要设置消息的 expiration 属性,该属性表示消息的过期时间。以下是发送延时消息的示例代码:

    Map headers = new HashMap<>();
    headers.put("x-delay", 5000);
    AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
            .headers(headers)
            .expiration("5000")
            .build();
    channel.basicPublish("delayed-exchange", "delayed-queue", properties, "Hello, delayed queue!".getBytes());

    在上述代码中,设置了消息的 expiration 属性为 5000 毫秒,并将消息发送到 delayed-exchange 交换机上,路由键为 delayed-queue,消息内容为 “Hello, delayed queue!”。

    消费延时消息

    消费延时消息时,需要设置消费者的 QOS(Quality of Service)参数,以控制消费者的并发处理能力。以下是消费延时消息的示例代码:

    channel.basicQos(1);
    channel.basicConsume("delayed-queue", false, (consumerTag, delivery) -> {
        String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
        System.out.println("Received message: " + message);
        channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
    });

    在上述代码中,设置了 QOS 参数为 1,即每次只处理一个消息。然后使用 basicConsume 方法消费 delayed-queue 队列中的消息,并在消费完成后,使用 basicAck 方法确认消息已被消费。

    通过上述步骤,就可以实现 RabbitMQ 延时队列,用于实现定时任务等功能。

    RabbitMQ延时队列是一种常见的消息队列应用场景,它可以在消息发送后指定一定的时间后才能被消费者消费,通常用于实现一些延时任务,例如订单超时未支付自动取消等。

    RabbitMQ延时队列具体代码

    下面是具体代码(附注释):

    import com.rabbitmq.client.*;
    import java.io.IOException;
    import java.util.HashMap;
    import java.util.Map;
    import java.util.concurrent.TimeoutException;
    
    public class DelayedQueueExample {
        private static final String EXCHANGE_NAME = "delayed_exchange";
        private static final String QUEUE_NAME = "delayed_queue";
        private static final String ROUTING_KEY = "delayed_routing_key";
    
        public static void main(String[] args) throws IOException, TimeoutException {
            ConnectionFactory factory = new ConnectionFactory();
            factory.setHost("localhost");
            Connection connection = factory.newConnection();
            Channel channel = connection.createChannel();
            /*
             Exchange.DeclareOk exchangeDeclare(String exchange,
                                                  String type,
                                                  boolean durable,
                                                  boolean autoDelete,
                                                  boolean internal,
                                                  Map arguments) throws IOException;
                                                  */
            // 创建一个支持延时队列的Exchange
            Map arguments = new HashMap<>();
            arguments.put("x-delayed-type", "direct");
            channel.exchangeDeclare(EXCHANGE_NAME, "x-delayed-message", true, false, arguments);
    
            // 创建一个延时队列,设置x-dead-letter-exchange和x-dead-letter-routing-key参数
            Map queueArguments = new HashMap<>();
            queueArguments.put("x-dead-letter-exchange", "");
            queueArguments.put("x-dead-letter-routing-key", QUEUE_NAME);
            queueArguments.put("x-message-ttl", 5000);
            channel.queueDeclare(QUEUE_NAME, true, false, false, queueArguments);
            channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
    
            // 发送消息到延时队列中,设置expiration参数
            AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
                    .expiration("10000")
                    .build();
            String message = "Hello, delayed queue!";
            channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, properties, message.getBytes());
            System.out.println("Sent message to delayed queue: " + message);
            channel.close();
            connection.close();
        }
    }

    在上面的代码中,我们创建了一个支持延时队列的Exchange,并创建了一个延时队列,设置了x-dead-letter-exchange和x-dead-letter-routing-key参数。然后,我们发送了一条消息到延时队列中,设置了expiration参数,表示这条消息延时10秒后才能被消费。

    注意,如果我们想要消费延时队列中的消息,需要创建一个消费者,并监听这个队列。当消息被消费时,需要发送ack确认消息已经被消费,否则消息会一直留在队列中。

    本文内容来源于互联网,如有侵权请联系删除。
    作者最新文章
    编程开发
    相关文章 更多
    谷歌浏览器Mac版入口
    谷歌浏览器Mac版入口

    谷歌浏览器Mac版官方安装指南 谷歌浏览器Mac版官方安装入口是https://www.google.com/chrome/,需macOS 12+系统、500MB空间,下载.dmg后拖入应用程序安装,支持多设备同步、性能优化与隐私保护功能。 苹果电脑Chrome的安装入口究竟在哪里?这个问题最近可是

    Chrome浏览器JS脚本不运行怎么办
    Chrome浏览器JS脚本不运行怎么办

    Chrome中JavaScript未执行需依次检查:一、移除站点级禁用并添加允许域名;二、开启全局JavaScript开关;三、禁用干扰扩展;四、在开发者工具中启用JavaScript;五、重置内容设置为默认。 有时在Chrome里打开网页,会发现交互按钮点了没反应,数据加载不出来,页面仿佛“静止”

    IE浏览器怀旧版在线网址
    IE浏览器怀旧版在线网址

    IE浏览器怀旧版在线网址:一次精准的技术时光回溯 最近,不少老用户和怀旧爱好者在反复搜索一个问题:那个经典的Internet Explorer,如今还能在哪里原汁原味地体验到?答案指向一个特定的地址:https://ie.microsoft.com/legacy/。 这个网站远不止是一个简单的“皮肤

    火狐浏览器有哪些设置功能
    火狐浏览器有哪些设置功能

    火狐浏览器五大核心设置功能:解锁高效、安全与个性化体验 火狐浏览器功能强大,但如果不仔细挖掘,很多能大幅提升效率和安全性的设置可能就“藏着掖着”了。这就好比拥有一台高性能设备,却只用了基础模式。那么,如何把它调整到最顺手、最安全的状态?接下来,我们就聚焦于当前版本(截至2025年末)最关键的五大设置

    chrome搜索免验证入口
    chrome搜索免验证入口

    Chrome官方免验证入口为https://www.google.cn/chrome/,提供全平台安装包、免登录即用、本地化安全机制及引擎级性能优化。 到底该去哪里找正版、免费且无需繁琐验证的Chrome浏览器入口?这个问题困扰了不少网友。今天,我们就来直通核心,为大家详细拆解Chrome引擎的官方

    java heap space 选型思路:使用场景与区别整理
    java heap space 选型思路:使用场景与区别整理

    Java堆是JVM存储对象的核心内存区域,配置需结合场景:单体应用适中设置;大数据处理需大堆并关注GC停顿;微服务强调快速启动;高并发需精细划分堆区域。关键参数-Xms和-Xmx建议等值以稳定性能。垃圾回收器选择影响效率,如G1适用于大堆,ZGC可实现低停顿。内存错误时需监控堆状态。

    java heap space 使用中遇到的问题怎么解决
    java heap space 使用中遇到的问题怎么解决

    Java堆内存溢出错误通常因内存泄漏、数据处理需求过大或JVM参数配置不当引起。排查时可借助jmap、堆转储及MAT等工具定位问题。解决方案包括调整JVM内存参数(如-Xmx)、修复代码中的内存泄漏、优化大数据处理逻辑,并建立持续监控与预防机制,以保障应用稳定运行。

    java xml 选型思路:使用场景与区别整理
    java xml 选型思路:使用场景与区别整理

    XML在Java开发中用于配置、数据交换等场景。解析方式主要有DOM、SAX、StAX及第三方库。DOM适合操作小文件,SAX/StAX适合处理大文件流,JAXB用于对象与XML映射。选型需结合数据大小、内存、性能及团队熟悉度,现代框架常封装底层解析。

    using namespace 使用中遇到的问题怎么解决
    using namespace 使用中遇到的问题怎么解决

    命名空间的基本概念与常见引入问题在C++等编程语言中,命名空间(namespace)是一种将代码标识符(如变量、函数、类名)封装在特定名称下的机制,其主要目的是避免命名冲突,尤其是在大型项目或使用多个第三方库时。使用“using namespace”指令可以将指定命名空间中的所有名称引入当前作用域,

    c语言函数递归 实操经验总结:这些技巧很实用
    c语言函数递归 实操经验总结:这些技巧很实用

    理解递归的基本原理在C语言中,递归是一种函数调用自身的编程技术。要掌握它,首先需要理解其核心思想:将一个复杂的大问题,分解为一个或几个与原问题相似但规模更小的子问题,直到子问题足够简单,可以直接求解。这个过程通常包含两个关键部分:递归出口和递归体。递归出口定义了问题何时不再继续分解,即最简单、可直接

    查看更多
    精品专题 更多
    装机必备
    装机必备

    正软商城装机必备专区,精选办公、浏览器、安全防护、影音播放、压缩解压、设计创作和系统工具等电脑常用正版软件,帮助用户快速完成新电脑软件配置。

    Windows
    Windows

    正软商城Windows软件专区,汇集适用于Windows电脑的办公、设计、安全防护、影音播放、开发工具和系统优化软件,提供软件介绍、系统要求、正版授权及购买下载服务。

    macOS软件
    macOS软件

    正软商城macOS软件专区,精选适用于Mac电脑的办公、设计、影音、效率、开发和系统工具,提供软件功能介绍、macOS兼容版本、正版授权及购买下载服务。

    Mac软件 更多
    灵活计算器
    灵活计算器
    macOS/iOS/Android

    灵活计算器是一款笔记式算数应用,支持实时计算、动态关联和云端同步功能。记录、整理和输出之间的过渡会更自然,适合长期写作、做笔记或持续沉淀个人内容。

    赤友清理大师
    赤友清理大师
    macOS

    赤友清理大师是一款为 Mac 设计的智能清理优化工具,可精准扫描垃圾、大文件、重复文件等,释放磁盘空间。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

    极度公式
    极度公式
    Windows/macOS/Linux

    极度公式是一款跨平台专业LaTeX公式识别编辑软件,支持OCR公式识别和多平台编辑。和使用说明,避免使用,享受完整功能与稳定支持。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

    WINDOWS 更多
    Windows 10
    Windows 10
    Windows

    Windows 10 是一款微软推出的经典操作系统,拥有硬件兼容性与多任务处理能力。它更偏向把系统状态查看和常用调节动作放在一起,适合需要持续观察和微调设备状态的场景。

    极度公式
    极度公式
    Windows/macOS/Linux

    极度公式是一款跨平台专业LaTeX公式识别编辑软件,支持OCR公式识别和多平台编辑。和使用说明,避免使用,享受完整功能与稳定支持。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

    密码键盘
    密码键盘
    Windows/macOS/iOS/Android

    密码键盘是一款兼具安全性与便捷性的高效密码管理器。日常使用里的持续防护和信息管理会更突出,适合把安全控制放进长期使用流程中的场景。