最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

使用Java實(shí)現(xiàn)RabbitMQ延時(shí)隊(duì)列

 更新時(shí)間:2023年06月20日 14:27:21   作者:土豆魚_  
RabbitMQ?延時(shí)隊(duì)列是指消息在發(fā)送到隊(duì)列后,并不立即被消費(fèi)者消費(fèi),而是等待一段時(shí)間后再被消費(fèi)者消費(fèi),本文為大家介紹了實(shí)現(xiàn)RabbitMQ延時(shí)隊(duì)列的Java代碼,希望對(duì)大家有所幫助

RabbitMQ 延時(shí)隊(duì)列介紹

RabbitMQ 延時(shí)隊(duì)列是指消息在發(fā)送到隊(duì)列后,并不立即被消費(fèi)者消費(fèi),而是等待一段時(shí)間后再被消費(fèi)者消費(fèi)。這種隊(duì)列通常用于實(shí)現(xiàn)定時(shí)任務(wù),例如,訂單超時(shí)未支付系統(tǒng)取消訂單釋放所占庫存等。

RabbitMQ實(shí)現(xiàn)延時(shí)隊(duì)列的方法有多種,其中比較常見的是使用插件或者通過DLX(Dead Letter Exchange)機(jī)制實(shí)現(xiàn)。

1.使用插件實(shí)現(xiàn)延時(shí)隊(duì)列

RabbitMQ提供了rabbitmq_delayed_message_exchange插件,可以通過該插件實(shí)現(xiàn)延時(shí)隊(duì)列。該插件的原理是在消息發(fā)送時(shí),將消息發(fā)送到一個(gè)特定的Exchange中,然后該Exchange會(huì)根據(jù)消息中的延時(shí)時(shí)間將消息轉(zhuǎn)發(fā)到指定的隊(duì)列中,從而實(shí)現(xiàn)延時(shí)隊(duì)列的功能

使用該插件需要先安裝插件,然后創(chuàng)建一個(gè)Exchange,并將該Exchange的類型設(shè)置為x-delayed-message,然后將該Exchange與隊(duì)列綁定即可。

2.使用DLX機(jī)制實(shí)現(xiàn)延時(shí)隊(duì)列

消息的TTL就是消息的存活時(shí)間。RabbitMQ可以對(duì)隊(duì)列和消息分別設(shè)置TTL。而對(duì)隊(duì)列設(shè)置就是隊(duì)列沒有消費(fèi)者連著的保留時(shí)間,也可以對(duì)每一個(gè)單獨(dú)的消息做單獨(dú)的 設(shè)置。超過了這個(gè)時(shí)間,我們認(rèn)為這個(gè)消息就死了,稱之為死信。如果隊(duì)列設(shè)置了,消息也設(shè)置了,那么會(huì)取小的。所以一個(gè)消息如果被路由到不同的隊(duì) 列中,這個(gè)消息死亡的時(shí)間有可能不一樣(不同的隊(duì)列設(shè)置)。這里單講單個(gè)消息的TTL,因?yàn)樗攀菍?shí)現(xiàn)延遲任務(wù)的關(guān)鍵。可以通過設(shè)置消息的expiration字段或者x- message-ttl屬性來設(shè)置時(shí)間,兩者是一樣的效果

DLX機(jī)制是RabbitMQ提供的一種消息轉(zhuǎn)發(fā)機(jī)制,它可以將無法被處理的消息轉(zhuǎn)發(fā)到指定的Exchange中,從而實(shí)現(xiàn)消息的延時(shí)處理。具體實(shí)現(xiàn)步驟如下:

  • 創(chuàng)建一個(gè)普通的Exchange和Queue,并將它們綁定在一起。
  • 創(chuàng)建一個(gè)DLX Exchange,并將普通Exchange綁定到該DLX Exchange上。
  • 將Queue設(shè)置為具有TTL(Time To Live)屬性,并設(shè)置消息過期時(shí)間。
  • 將Queue綁定到DLX Exchange上。

當(dāng)消息過期后,會(huì)被發(fā)送到DLX Exchange中,然后再由DLX Exchange將消息轉(zhuǎn)發(fā)到指定的Exchange中,從而實(shí)現(xiàn)延時(shí)隊(duì)列的功能。

使用DLX機(jī)制實(shí)現(xiàn)延時(shí)隊(duì)列的優(yōu)點(diǎn)是不需要安裝額外的插件,但是需要對(duì)消息的過期時(shí)間進(jìn)行精確控制,否則可能會(huì)出現(xiàn)消息過期時(shí)間不準(zhǔn)確的情況。

Java語言設(shè)置延時(shí)隊(duì)列

下面是使用 Java 語言通過 RabbitMQ 設(shè)置延時(shí)隊(duì)列的步驟:

1.安裝插件

首先,需要安裝 rabbitmq_delayed_message_exchange 插件??梢酝ㄟ^以下命令安裝:

rabbitmq-plugins enable rabbitmq_delayed_message_exchange

2. 創(chuàng)建延時(shí)交換機(jī)

延時(shí)隊(duì)列需要使用延時(shí)交換機(jī)??梢允褂?x-delayed-message 類型創(chuàng)建一個(gè)延時(shí)交換機(jī)。以下是創(chuàng)建延時(shí)交換機(jī)的示例代碼:

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

3.創(chuàng)建延時(shí)隊(duì)列

創(chuàng)建延時(shí)隊(duì)列時(shí),需要將隊(duì)列綁定到延時(shí)交換機(jī)上,并設(shè)置隊(duì)列的 TTL(Time To Live)參數(shù)。以下是創(chuàng)建延時(shí)隊(duì)列的示例代碼:

Map<String, Object> 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");

在上述代碼中,將隊(duì)列綁定到延時(shí)交換機(jī)上,并設(shè)置了隊(duì)列的 TTL 參數(shù)為 5000 毫秒,即消息在發(fā)送到隊(duì)列后,如果在 5000 毫秒內(nèi)沒有被消費(fèi)者消費(fèi),則會(huì)被轉(zhuǎn)發(fā)到 delayed-exchange 交換機(jī)上,并發(fā)送到 delayed-queue 隊(duì)列中。

4.發(fā)送延時(shí)消息

發(fā)送延時(shí)消息時(shí),需要設(shè)置消息的 expiration 屬性,該屬性表示消息的過期時(shí)間。以下是發(fā)送延時(shí)消息的示例代碼:

Map<String, Object> 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());

在上述代碼中,設(shè)置了消息的 expiration 屬性為 5000 毫秒,并將消息發(fā)送到 delayed-exchange 交換機(jī)上,路由鍵為 delayed-queue,消息內(nèi)容為 "Hello, delayed queue!"。

5.消費(fèi)延時(shí)消息

消費(fèi)延時(shí)消息時(shí),需要設(shè)置消費(fèi)者的 QOS(Quality of Service)參數(shù),以控制消費(fèi)者的并發(fā)處理能力。以下是消費(fèi)延時(shí)消息的示例代碼:

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);
});

在上述代碼中,設(shè)置了 QOS 參數(shù)為 1,即每次只處理一個(gè)消息。然后使用 basicConsume 方法消費(fèi) delayed-queue 隊(duì)列中的消息,并在消費(fèi)完成后,使用 basicAck 方法確認(rèn)消息已被消費(fèi)。

通過上述步驟,就可以實(shí)現(xiàn) RabbitMQ 延時(shí)隊(duì)列,用于實(shí)現(xiàn)定時(shí)任務(wù)等功能。

RabbitMQ延時(shí)隊(duì)列是一種常見的消息隊(duì)列應(yīng)用場景,它可以在消息發(fā)送后指定一定的時(shí)間后才能被消費(fèi)者消費(fèi),通常用于實(shí)現(xiàn)一些延時(shí)任務(wù),例如訂單超時(shí)未支付自動(dòng)取消等。

RabbitMQ延時(shí)隊(duì)列具體代碼

下面是具體代碼(附注釋):

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<String, Object> arguments) throws IOException;
                                              */
        // 創(chuàng)建一個(gè)支持延時(shí)隊(duì)列的Exchange
        Map<String, Object> arguments = new HashMap<>();
        arguments.put("x-delayed-type", "direct");
        channel.exchangeDeclare(EXCHANGE_NAME, "x-delayed-message", true, false, arguments);
        // 創(chuàng)建一個(gè)延時(shí)隊(duì)列,設(shè)置x-dead-letter-exchange和x-dead-letter-routing-key參數(shù)
        Map<String, Object> 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);
        // 發(fā)送消息到延時(shí)隊(duì)列中,設(shè)置expiration參數(shù)
        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();
    }
}

在上面的代碼中,我們創(chuàng)建了一個(gè)支持延時(shí)隊(duì)列的Exchange,并創(chuàng)建了一個(gè)延時(shí)隊(duì)列,設(shè)置了x-dead-letter-exchange和x-dead-letter-routing-key參數(shù)。然后,我們發(fā)送了一條消息到延時(shí)隊(duì)列中,設(shè)置了expiration參數(shù),表示這條消息延時(shí)10秒后才能被消費(fèi)。

注意,如果我們想要消費(fèi)延時(shí)隊(duì)列中的消息,需要?jiǎng)?chuàng)建一個(gè)消費(fèi)者,并監(jiān)聽這個(gè)隊(duì)列。當(dāng)消息被消費(fèi)時(shí),需要發(fā)送ack確認(rèn)消息已經(jīng)被消費(fèi),否則消息會(huì)一直留在隊(duì)列中。

到此這篇關(guān)于使用Java實(shí)現(xiàn)RabbitMQ延時(shí)隊(duì)列的文章就介紹到這了,更多相關(guān)Java RabbitMQ延時(shí)隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java語言中4種內(nèi)部類的超詳細(xì)講解

    Java語言中4種內(nèi)部類的超詳細(xì)講解

    這篇文章主要給大家介紹了關(guān)于Java語言中4種內(nèi)部類的超詳細(xì)講解,內(nèi)部類可以分為:實(shí)例內(nèi)部類、靜態(tài)內(nèi)部類和成員內(nèi)部類,每種內(nèi)部類都有它特定的一些特點(diǎn),文中介紹的非常詳細(xì),需要的朋友可以參考下
    2023-04-04
  • 詳解Java中的ThreadLocal

    詳解Java中的ThreadLocal

    ThreadLocal是JDK包提供的,它提供線程本地變量,如果創(chuàng)建一個(gè)ThreadLocal變量,那么訪問這個(gè)變量的每個(gè)線程都會(huì)有這個(gè)變量的一個(gè)副本,在實(shí)際多線程操作的時(shí)候,操作的是自己本地內(nèi)存中的變量,從而規(guī)避了線程安全問題
    2021-06-06
  • Java多線程并發(fā)編程(互斥鎖Reentrant Lock)

    Java多線程并發(fā)編程(互斥鎖Reentrant Lock)

    這篇文章主要介紹了ReentrantLock 互斥鎖,在同一時(shí)間只能被一個(gè)線程所占有,在被持有后并未釋放之前,其他線程若想獲得該鎖只能等待或放棄,需要的朋友可以參考下
    2017-05-05
  • 使用idea和gradle編譯spring5源碼的方法步驟

    使用idea和gradle編譯spring5源碼的方法步驟

    這篇文章主要介紹了詳解使用idea和gradle編譯spring5源碼,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-12-12
  • Java響應(yīng)式編程之Flux與SseEmitter深度解析(附詳細(xì)代碼)

    Java響應(yīng)式編程之Flux與SseEmitter深度解析(附詳細(xì)代碼)

    這篇文章主要介紹了Java響應(yīng)式編程之Flux與SseEmitter深度解析的相關(guān)資料,從需求分析、技術(shù)背景、使用方法、底層原理、性能對(duì)比到生產(chǎn)環(huán)境實(shí)戰(zhàn),全面解析了這兩種技術(shù)的特點(diǎn)、應(yīng)用場景及優(yōu)缺點(diǎn),需要的朋友可以參考下
    2026-01-01
  • Java合并區(qū)間的實(shí)現(xiàn)

    Java合并區(qū)間的實(shí)現(xiàn)

    本文主要介紹了Java合并區(qū)間的實(shí)現(xiàn),通過合理使用集合類和排序算法,可以有效地解決合并區(qū)間問題,具有一定的參考價(jià)值,感興趣的可以了解一下
    2023-08-08
  • Java長度不足左位補(bǔ)0的3種實(shí)現(xiàn)方法

    Java長度不足左位補(bǔ)0的3種實(shí)現(xiàn)方法

    這篇文章主要介紹了Java長度不足左位補(bǔ)0的3種實(shí)現(xiàn)方法小結(jié),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-12-12
  • java生成在線驗(yàn)證碼

    java生成在線驗(yàn)證碼

    這篇文章主要介紹了java生成在線驗(yàn)證碼,需要的朋友可以參考下
    2023-10-10
  • JDK1.8使用的垃圾回收器和執(zhí)行GC的時(shí)長以及GC的頻率方式

    JDK1.8使用的垃圾回收器和執(zhí)行GC的時(shí)長以及GC的頻率方式

    這篇文章主要介紹了JDK1.8使用的垃圾回收器和執(zhí)行GC的時(shí)長以及GC的頻率方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-05-05
  • SpringBoot 圖書管理系統(tǒng)(刪除、強(qiáng)制登錄、更新圖書)詳細(xì)代碼

    SpringBoot 圖書管理系統(tǒng)(刪除、強(qiáng)制登錄、更新圖書)詳細(xì)代碼

    在企業(yè)開發(fā)中,通常不采用delete語句進(jìn)行物理刪除,而是使用邏輯刪除,邏輯刪除通過修改標(biāo)識(shí)字段來表示數(shù)據(jù)已被刪除,方便數(shù)據(jù)恢復(fù),本文給大家介紹SpringBoot 圖書管理系統(tǒng)實(shí)例代碼,感興趣的朋友跟隨小編一起看看吧
    2024-09-09

最新評(píng)論

平潭县| 通榆县| 烟台市| 万全县| 白城市| 察隅县| 剑川县| 楚雄市| 蓬莱市| 凤庆县| 四子王旗| 芜湖市| 永康市| 仁布县| 溧阳市| 贵州省| 利津县| 大厂| 桐庐县| 福泉市| 木里| 天镇县| 新蔡县| 旺苍县| 斗六市| 建瓯市| 通化县| 项城市| 汤阴县| 水富县| 大兴区| 安福县| 宁明县| 麻城市| 黔江区| 贵州省| 临城县| 丹江口市| 田东县| 建始县| 托克逊县|