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

Java?RabbitMQ的持久化和發(fā)布確認(rèn)詳解

 更新時(shí)間:2022年03月08日 11:49:06   作者:江海i  
這篇文章主要為大家詳細(xì)介紹了RabbitMQ的持久化和發(fā)布確認(rèn),文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下,希望能夠給你帶來(lái)幫助

1. 持久化

當(dāng)RabbitMQ服務(wù)停掉以后消息生產(chǎn)者發(fā)送過(guò)的消息不丟失。默認(rèn)情況下RabbitMQ退出或者崩潰時(shí),會(huì)忽視掉隊(duì)列和消息。為了保證消息不丟失需要將隊(duì)列和消息都標(biāo)記為持久化。

1.1 實(shí)現(xiàn)持久化

1.隊(duì)列持久化:在創(chuàng)建隊(duì)列時(shí)將channel.queueDeclare();第二個(gè)參數(shù)改為true。

2.消息持久化:在使用信道發(fā)送消息時(shí)channel.basicPublish();將第三個(gè)參數(shù)改為:MessageProperties.PERSISTENT_TEXT_PLAIN表示持久化消息。

/**
 * @Description 持久化MQ
 * @date 2022/3/7 9:14
 */
public class Producer3 {
    private static final String LONG_QUEUE = "long_queue";
    public static void main(String[] args) throws Exception {
        Channel channel = RabbitMQUtils.getChannel();
        // 持久化隊(duì)列
        channel.queueDeclare(LONG_QUEUE,true,false,false,null);
        Scanner scanner = new Scanner(System.in);
        int i = 0;
        while (scanner.hasNext()){
            i++;
            String msg = scanner.next() + i;
            // 持久化消息
            channel.basicPublish("",LONG_QUEUE, MessageProperties.PERSISTENT_TEXT_PLAIN,msg.getBytes(StandardCharsets.UTF_8));
            System.out.println("發(fā)送消息:'" + msg + "'成功");
        }
    }
}

但是存儲(chǔ)消息還有存在一個(gè)緩存的間隔點(diǎn),沒(méi)有真正的寫(xiě)入磁盤(pán),持久性保證不夠強(qiáng),但是對(duì)于簡(jiǎn)單隊(duì)列而言也綽綽有余。

1.2 不公平分發(fā)

輪詢分發(fā)的方式在消費(fèi)者處理效率不同的情況下并不適用。所以真正的公平應(yīng)該是遵循能者多勞的前提。

在消費(fèi)者處修改channel.basicQos(1);表示開(kāi)啟不公平分發(fā)

/**
 * @Description 不公平分發(fā)消費(fèi)者
 * @date 2022/3/7 9:27
 */
public class Consumer2 {
    private static final String LONG_QUEUE = "long_queue";
    public static void main(String[] args) throws Exception {
        Channel channel = RabbitMQUtils.getChannel();
        DeliverCallback deliverCallback = (consumerTag, message) -> {
            // 模擬并發(fā)沉睡三十秒
            try {
                Thread.sleep(30000);
                System.out.println("線程B接收消息:"+ new String(message.getBody(), StandardCharsets.UTF_8));
                channel.basicAck(message.getEnvelope().getDeliveryTag(),false);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        };
        // 設(shè)置不公平分發(fā)
        channel.basicQos(1);
        channel.basicConsume(LONG_QUEUE,false,deliverCallback,
                consumerTag -> {
                    System.out.println(consumerTag + "消費(fèi)者取消消費(fèi)");
                });
    }
}

1.3 測(cè)試不公平分發(fā)

測(cè)試目的:是否能實(shí)現(xiàn)能者多勞。

測(cè)試方法:兩個(gè)消費(fèi)者睡眠不同的事件來(lái)模擬處理事件不同,如果處理時(shí)間(睡眠時(shí)間)短的能夠處理多個(gè)消息就代表目的達(dá)成。

先啟動(dòng)生產(chǎn)者創(chuàng)建隊(duì)列,再分別啟動(dòng)兩個(gè)消費(fèi)者。

生產(chǎn)者按照順序發(fā)四條消息:

在這里插入圖片描述

睡眠時(shí)間短的線程A接收到了三條消息

在這里插入圖片描述

而睡眠時(shí)間長(zhǎng)的線程B只接收到的第二條消息:

在這里插入圖片描述

因?yàn)榫€程B在處理消息時(shí)消耗的時(shí)間較長(zhǎng),所以就將其他消息分配給了線程A。

實(shí)驗(yàn)成功!

1.4 預(yù)取值

消息的發(fā)送和手動(dòng)確認(rèn)都是異步完成的,因此就存在一個(gè)未確認(rèn)消息的緩沖區(qū),開(kāi)發(fā)人員希望能夠限制緩沖區(qū)的大小,用來(lái)避免緩沖區(qū)里面無(wú)限制的未確認(rèn)消息問(wèn)題。

這里的預(yù)期值就值得是上述方法channel.basicQos();里面的參數(shù),如果在當(dāng)前信道上存在等于參數(shù)的消息就不會(huì)在安排當(dāng)前信道進(jìn)行消費(fèi)消息。

1.4.1 代碼測(cè)試

測(cè)試方法:

1.新建兩個(gè)不同的消費(fèi)者分別給定預(yù)期值5個(gè)2。

2.給睡眠時(shí)間長(zhǎng)的指定為5,時(shí)間短的指定為2。

3.假如按照指定的預(yù)期值獲取消息則表示測(cè)試成功,但并不是代表一定會(huì)按照5和2分配,這個(gè)類(lèi)似于權(quán)重的判別。

代碼根據(jù)上述代碼修改預(yù)期值即可。

2. 發(fā)布確認(rèn)

發(fā)布確認(rèn)就是生產(chǎn)者發(fā)布消息到隊(duì)列之后,隊(duì)列確認(rèn)進(jìn)行持久化完畢再通知給生產(chǎn)者的過(guò)程。這樣才能保證消息不會(huì)丟失。

需要注意的是需要開(kāi)啟隊(duì)列持久化才能使用確認(rèn)發(fā)布。
開(kāi)啟方法:channel.confirmSelect();

2.1 單個(gè)確認(rèn)發(fā)布

是一種同步發(fā)布的方式,即發(fā)送完一個(gè)消息之后只有確認(rèn)它確認(rèn)發(fā)布后,后續(xù)的消息才會(huì)繼續(xù)發(fā)布,在指定的時(shí)間內(nèi)沒(méi)有確認(rèn)就會(huì)拋出異常。缺點(diǎn)就是特別慢。

/**
 * @Description 確認(rèn)發(fā)布——單個(gè)確認(rèn)
 * @date 2022/3/7 14:49
 */
public class SoloProducer {
    private static final int MESSAGE_COUNT = 100;
    private static final String QUEUE_NAME = "confirm_solo";
    public static void main(String[] args) throws Exception {
        Channel channel = RabbitMQUtils.getChannel();
        // 產(chǎn)生隊(duì)列
        channel.queueDeclare(QUEUE_NAME,true,false,false,null);
        // 開(kāi)啟確認(rèn)發(fā)布
        channel.confirmSelect();
        // 記錄開(kāi)始時(shí)間
        long beginTime = System.currentTimeMillis();
        for (int i = 0; i < MESSAGE_COUNT; i++) {
            String msg = ""+i;
            channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,msg.getBytes(StandardCharsets.UTF_8));
            // 單個(gè)發(fā)布確認(rèn)
            boolean flag = channel.waitForConfirms();
            if (flag){
                System.out.println("發(fā)送消息:" + i);
            }
        }
        // 記錄結(jié)束時(shí)間
        long endTime = System.currentTimeMillis();
        System.out.println("發(fā)送" + MESSAGE_COUNT + "條消息消耗:"+(endTime - beginTime) + "毫秒");   }
}

2.2 批量確認(rèn)發(fā)布

一批一批的確認(rèn)發(fā)布可以提高系統(tǒng)的吞吐量。但是缺點(diǎn)是發(fā)生故障導(dǎo)致發(fā)布出現(xiàn)問(wèn)題時(shí),需要將整個(gè)批處理保存在內(nèi)存中,后面再重新發(fā)布。

/**
 * @Description 確認(rèn)發(fā)布——批量確認(rèn)
 * @date 2022/3/7 14:49
 */
public class BatchProducer {
    private static final int MESSAGE_COUNT = 100;
    private static final String QUEUE_NAME = "confirm_batch";
    public static void main(String[] args) throws Exception {
        Channel channel = RabbitMQUtils.getChannel();
        // 產(chǎn)生隊(duì)列
        channel.queueDeclare(QUEUE_NAME,true,false,false,null);
        // 開(kāi)啟確認(rèn)發(fā)布
        channel.confirmSelect();
        // 設(shè)置一個(gè)多少一批確認(rèn)一次。
        int batchSize = MESSAGE_COUNT / 10;
        // 記錄開(kāi)始時(shí)間
        long beginTime = System.currentTimeMillis();
        for (int i = 0; i < MESSAGE_COUNT; i++) {
            String msg = ""+i;
            channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,msg.getBytes(StandardCharsets.UTF_8));
            // 批量發(fā)布確認(rèn)
            if (i % batchSize == 0){
                if (channel.waitForConfirms()){
                    System.out.println("發(fā)送消息:" + i);
                }
            }
        }
        // 記錄結(jié)束時(shí)間
        long endTime = System.currentTimeMillis();
        System.out.println("發(fā)送" + MESSAGE_COUNT + "條消息消耗:"+(endTime - beginTime) + "毫秒");
    }
}

顯然效率要比單個(gè)確認(rèn)發(fā)布的高很多。

2.3 異步確認(rèn)發(fā)布

在編程上比上述兩個(gè)要復(fù)雜,但是性價(jià)比很高,無(wú)論是可靠性還行效率的都好很多,利用回調(diào)函數(shù)來(lái)達(dá)到消息可靠性傳遞的。

/**
 * @Description 確認(rèn)發(fā)布——異步確認(rèn)
 * @date 2022/3/7 14:49
 */
public class AsyncProducer {
    private static final int MESSAGE_COUNT = 100;
    private static final String QUEUE_NAME = "confirm_async";
    public static void main(String[] args) throws Exception {
        Channel channel = RabbitMQUtils.getChannel();
        // 產(chǎn)生隊(duì)列
        channel.queueDeclare(QUEUE_NAME,true,false,false,null);
        // 開(kāi)啟確認(rèn)發(fā)布
        channel.confirmSelect();
        // 記錄開(kāi)始時(shí)間
        long beginTime = System.currentTimeMillis();
        // 確認(rèn)成功回調(diào)
        ConfirmCallback ackCallback = (deliveryTab,multiple) ->{
            System.out.println("確認(rèn)成功消息:" + deliveryTab);
        };
        // 確認(rèn)失敗回調(diào)
        ConfirmCallback nackCallback = (deliveryTab,multiple) ->{
            System.out.println("未確認(rèn)的消息:" + deliveryTab);
        };
        // 消息監(jiān)聽(tīng)器
        /**
         * addConfirmListener:
         *                  1. 確認(rèn)成功的消息;
         *                  2. 確認(rèn)失敗的消息。
         */
        channel.addConfirmListener(ackCallback,nackCallback);
        for (int i = 0; i < MESSAGE_COUNT; i++) {
            String msg = "" + i;
            channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,msg.getBytes(StandardCharsets.UTF_8));
        }

        // 記錄結(jié)束時(shí)間
        long endTime = System.currentTimeMillis();
        System.out.println("發(fā)送" + MESSAGE_COUNT + "條消息消耗:"+(endTime - beginTime) + "毫秒");
    }
}

2.4 處理未確認(rèn)的消息

最好的處理方式把未確認(rèn)的消息放到一個(gè)基于內(nèi)存的能被發(fā)布線程訪問(wèn)的隊(duì)列。

例如:ConcurrentLinkedQueue可以在確認(rèn)隊(duì)列confirm callbacks與發(fā)布線程之間進(jìn)行消息的傳遞。

處理方式:

1.記錄要發(fā)送的全部消息;

2.在發(fā)布成功確認(rèn)處刪除;

3.打印未確認(rèn)的消息。

使用一個(gè)哈希表存儲(chǔ)消息,它的優(yōu)點(diǎn):

可以將需要和消息進(jìn)行關(guān)聯(lián);輕松批量刪除條目;支持高并發(fā)。

ConcurrentSkipListMap<Long,String > map = new ConcurrentSkipListMap<>();
/**
 * @Description 異步發(fā)布確認(rèn),處理未發(fā)布成功的消息
 * @date 2022/3/7 18:09
 */
public class AsyncProducerRemember {
    private static final int MESSAGE_COUNT = 100;
    private static final String QUEUE_NAME = "confirm_async_remember";
    public static void main(String[] args) throws Exception {
        Channel channel = RabbitMQUtils.getChannel();
        // 產(chǎn)生隊(duì)列
        channel.queueDeclare(QUEUE_NAME,true,false,false,null);
        // 開(kāi)啟確認(rèn)發(fā)布
        channel.confirmSelect();
        // 線程安全有序的一個(gè)hash表,適用與高并發(fā)
        ConcurrentSkipListMap< Long, String > map = new ConcurrentSkipListMap<>();
        // 記錄開(kāi)始時(shí)間
        long beginTime = System.currentTimeMillis();
        // 確認(rèn)成功回調(diào)
        ConfirmCallback ackCallback = (deliveryTab, multiple) ->{
            //2. 在發(fā)布成功確認(rèn)處刪除;
            // 批量刪除
            if (multiple){
                ConcurrentNavigableMap<Long, String> confirmMap = map.headMap(deliveryTab);
                confirmMap.clear();
            }else {
                // 單獨(dú)刪除
                map.remove(deliveryTab);
            }
            System.out.println("確認(rèn)成功消息:" + deliveryTab);
        };
        // 確認(rèn)失敗回調(diào)
        ConfirmCallback nackCallback = (deliveryTab,multiple) ->{
            // 3. 打印未確認(rèn)的消息。
            System.out.println("未確認(rèn)的消息:" + map.get(deliveryTab) + ",標(biāo)記:" + deliveryTab);
        };
        // 消息監(jiān)聽(tīng)器
        /**
         * addConfirmListener:
         *                  1. 確認(rèn)成功的消息;
         *                  2. 確認(rèn)失敗的消息。
         */
        channel.addConfirmListener(ackCallback,nackCallback);
        for (int i = 0; i < MESSAGE_COUNT; i++) {
            String msg = "" + i;
            channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,msg.getBytes(StandardCharsets.UTF_8));
            // 1. 記錄要發(fā)送的全部消息;
            map.put(channel.getNextPublishSeqNo(),msg);
        }

        // 記錄結(jié)束時(shí)間
        long endTime = System.currentTimeMillis();
        System.out.println("發(fā)送" + MESSAGE_COUNT + "條消息消耗:"+(endTime - beginTime) + "毫秒");
    }
}

總結(jié)

顯然來(lái)說(shuō),異步處理除了在編碼處有些麻煩,在處理時(shí)間效率和可用性上都是比單處理和批處理好很多。

本篇文章就到這里了,希望能夠給你帶來(lái)幫助,也希望您能夠多多關(guān)注腳本之家的更多內(nèi)容!    

相關(guān)文章

  • Spring事務(wù)執(zhí)行流程及如何創(chuàng)建事務(wù)

    Spring事務(wù)執(zhí)行流程及如何創(chuàng)建事務(wù)

    這篇文章主要介紹了Spring事務(wù)執(zhí)行流程及如何創(chuàng)建事務(wù),幫助大家更好的理解和學(xué)習(xí)使用spring框架,感興趣的朋友可以了解下
    2021-03-03
  • java中for循環(huán)執(zhí)行的順序圖文詳析

    java中for循環(huán)執(zhí)行的順序圖文詳析

    關(guān)于java的for循環(huán)想必大家非常熟悉,它是java常用的語(yǔ)句之一,這篇文章主要給大家介紹了關(guān)于java中for循環(huán)執(zhí)行順序的相關(guān)資料,需要的朋友可以參考下
    2021-06-06
  • SpringBoot熱部署啟動(dòng)關(guān)閉流程詳解

    SpringBoot熱部署啟動(dòng)關(guān)閉流程詳解

    Spring?Boot啟動(dòng)熱部署是一種技術(shù),它能讓開(kāi)發(fā)者在不重啟應(yīng)用程序的情況下實(shí)時(shí)更新代碼。這樣可以提高開(kāi)發(fā)效率,避免頻繁重啟應(yīng)用,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)吧
    2023-04-04
  • Java模板方法模式定義算法框架

    Java模板方法模式定義算法框架

    Java模板方法模式是一種行為型設(shè)計(jì)模式,它定義了一個(gè)算法框架,由抽象父類(lèi)定義算法的基本結(jié)構(gòu),具體實(shí)現(xiàn)細(xì)節(jié)由子類(lèi)來(lái)實(shí)現(xiàn),從而實(shí)現(xiàn)代碼復(fù)用和擴(kuò)展性
    2023-05-05
  • Java-URLDecoder、URLEncoder使用及說(shuō)明

    Java-URLDecoder、URLEncoder使用及說(shuō)明

    本文介紹了Java中URLDecoder和URLEncoder類(lèi)的使用方法,包括編碼和解碼規(guī)則、推薦的編碼方案、解碼器處理非法字符的方法以及URL編碼和解碼的示例
    2024-12-12
  • 詳解Spring mvc DispatchServlet 實(shí)現(xiàn)機(jī)制

    詳解Spring mvc DispatchServlet 實(shí)現(xiàn)機(jī)制

    本篇文章主要介紹了詳解Spring mvc DispatchServlet 實(shí)現(xiàn)機(jī)制,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2017-09-09
  • TransmittableThreadLocal線程間傳遞邏輯示例解析

    TransmittableThreadLocal線程間傳遞邏輯示例解析

    這篇文章主要介紹了TransmittableThreadLocal線程間傳遞邏輯示例解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-06-06
  • Java發(fā)送post方法詳解

    Java發(fā)送post方法詳解

    這篇文章主要介紹了Java發(fā)送post方法,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2019-04-04
  • Java匿名對(duì)象與匿名內(nèi)部類(lèi)

    Java匿名對(duì)象與匿名內(nèi)部類(lèi)

    本篇文章給大家詳細(xì)講解了Java匿名對(duì)象與匿名內(nèi)部類(lèi)的相關(guān)知識(shí)點(diǎn),有興趣的讀者一起學(xué)習(xí)下。
    2018-03-03
  • 詳解Spring中的@Scope注解

    詳解Spring中的@Scope注解

    這篇文章主要介紹了詳解Spring中的@Scope注解,@Scope注解是Spring IOC容器中的一個(gè)作用域,在Spring IOC容器中,他用來(lái)配置Bean實(shí)例的作用域?qū)ο?需要的朋友可以參考下
    2023-07-07

最新評(píng)論

永嘉县| 惠来县| 西昌市| 广河县| 湘潭县| 伽师县| 瑞昌市| 治多县| 舒兰市| 苍南县| 厦门市| 沾益县| 铜川市| 沿河| 南和县| 辽宁省| 马鞍山市| 古丈县| 南康市| 晴隆县| 西昌市| 化州市| 乐昌市| 双鸭山市| 光泽县| 衡南县| 明溪县| 尼勒克县| 嵊州市| 银川市| 清新县| 大同县| 琼海市| 乌兰县| 岳阳市| 彭阳县| 方正县| 沅陵县| 汾西县| 苍溪县| 灯塔市|