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

詳解RabbitMq如何做到消息的可靠性投遞

 更新時間:2022年09月05日 14:37:05   作者:S1C  
這篇文章主要為大家介紹了RabbitMq如何做到消息的可靠性投遞,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

前言

現(xiàn)在的一些互聯(lián)網(wǎng)項目或者是高并發(fā)的項目中很少有沒有引入消息隊列的。 引入消息隊列可以給這個項目帶來很多的好處:比如

  • 削峰

這個就很好的理解,在系統(tǒng)中的請求量是固定的,但是有的時候會多出很多的突發(fā)流量,比如在有秒殺活動的時候,這種瞬時的高流量可能會打垮系統(tǒng),這個時候就可以很好的引入MQ,將這些請求積壓到MQ中,然后消費端在按照自已的能力去處理這里請求

  • 解耦合

比如現(xiàn)在有系統(tǒng)A,當系統(tǒng)A執(zhí)行完成后,B、C系統(tǒng)需要拿到A系統(tǒng)的結(jié)果才可以繼續(xù)執(zhí)行,如果不引入MQ,A系統(tǒng)還要調(diào)用B、C系統(tǒng),這樣這A、B、C三個系統(tǒng)的耦合性就很大。引入MQ后A系統(tǒng)的執(zhí)行結(jié)果只需要保證將消息投遞到MQ就好,其它的兩個系統(tǒng)只需要監(jiān)聽這個MQ的某個隊列,這樣就降低了這三個系統(tǒng)之間的耦合性。

  • 異步

再通過A、B、C這三個系統(tǒng)舉例。A系統(tǒng)在返回給用戶的執(zhí)行結(jié)果前需要完成B、C系統(tǒng)的調(diào)用,這個總的執(zhí)行時間是A+B+C的執(zhí)行時間,如果引入MQ,A系統(tǒng)的執(zhí)行完成后將數(shù)據(jù)投遞到MQ,直接響應(yīng)用戶。B、C再這在通過監(jiān)聽完成數(shù)據(jù)的處理。這樣也降低了用戶的等待時間

除了這些好處,當然引入MQ還會有不好的地方:比如

  • 數(shù)據(jù)一致性問題
    • A系統(tǒng)執(zhí)行完將數(shù)據(jù)投遞到了MQ,B、C在消費的時候如果出現(xiàn)了問題,是不是就導(dǎo)致了數(shù)據(jù)不一致的問題
  • 可用性降低
    • 一個好好的系統(tǒng),引入一個MQ,如果這個MQ拓機了呢?這個可能就需要集群來提高MQ的高可用。
  • 系統(tǒng)的復(fù)雜度提高
    • 引入了MQ,我們還需要關(guān)注消息是否被成功的投遞,MQ中的消息被積壓太多怎么辦?消費端是否成功的消費的消息。

這些都是問題,所在是否要引入MQ還需要看業(yè)務(wù)需求

RabbitMq的投遞及消費流程

這里有張投遞消息到消費的流程圖

從這張圖上可看出這也是一種AMQP協(xié)議的實現(xiàn)。消息的提供者先是通過某一個信道將消息發(fā)送到交換機,然后交換機通通RoutingKey來將消息分發(fā)到某一個隊列上。然后,消費者在臨聽某一個隊列來進行消息的消費。

今天我們的主題是如何保證消息的投遞可靠性。那么我們來想想在這個流程中那些位置可能會影響我們消息的投遞可靠性?

從上圖中我們可以總結(jié)出有二個因素影響著消息是否被成功投遞和被成功消費

提供者

  • 提供者有沒有將消成功的發(fā)送到MQ并被處理
  • 發(fā)送到MQ中的消息有沒有成功的被路由到隊列中

消費者

  • 消費者有沒有成功的簽收消息并成功處理。
  • 消費者是否可以保證消費者的穩(wěn)定性

提供者如何確保消息的成功投遞

解決這個問題,我們可以通過提供者的發(fā)送方確認機制來實現(xiàn),這個發(fā)送方確認機制又分成三種:

  • 單條消息的同步確認
  • 多條消息的同步確認
  • 異步消息確認

單條消息的同步確認

首先要在當前的Channel上開啟消息確認模式,然后通過waitForConfirms()方法進行消息確認是否發(fā)送成功。

public static void main(String[] args) throws InterruptedException, TimeoutException, IOException {
        ConnectionFactory cf = new ConnectionFactory();
        cf.setHost("host");
        cf.setPort(5672);
        cf.setUsername("賬號");
        cf.setPassword("密碼");
        try(Connection connection = cf.newConnection();
            Channel channel = connection.createChannel()){
            channel.confirmSelect();
            Map<String,String> mes = new HashMap<>();
            mes.put("name","1111");
            String messageStr = objectMapper.writeValueAsString(mes);
            channel.basicPublish(
                    "exchange.drinks",
                    "drinks.juzi",
                    null,
                    messageStr.getBytes());
            boolean isSendSuccess = channel.waitForConfirms();
            if(isSendSuccess){
                System.out.print("消息發(fā)送成功");
            }
        }
    }

這樣做的話每次發(fā)完消息后,都會確保消息是否發(fā)送成功。如果發(fā)送失敗的話進行相應(yīng)的處理。

多條消息的同步確認

多條消息的確認和單條的差不多,比如我將發(fā)送消息的代碼放到一個循環(huán)內(nèi)。

public static void main(String[] args) throws InterruptedException, TimeoutException, IOException {
        ConnectionFactory cf = new ConnectionFactory();
        cf.setHost("host");
        cf.setPort(5672);
        cf.setUsername("賬號");
        cf.setPassword("密碼");
        try(Connection connection = cf.newConnection();
            Channel channel = connection.createChannel()){
            channel.confirmSelect();
            Map<String,String> mes = new HashMap<>();
            mes.put("name","1111");
            String messageStr = objectMapper.writeValueAsString(mes);
            for(int i = 0;i < 100;i++){
                channel.basicPublish(
                        "exchange.drinks",
                        "drinks.juzi",
                        null,
                        messageStr());
            }
            boolean isSendSuccess = channel.waitForConfirms();
            System.out.println(isSendSuccess);
        }
    }

這樣的話當一批消息發(fā)送完成后,進行統(tǒng)一的消息確認是否發(fā)送成功,就成了多條的消息確認,不過并不推薦使用這種確認消息的方式

在多條的消息確認中,比如我先是發(fā)送了一批的消息,比如這批消息有100條,這個時候如果有其中的一條消息沒有發(fā)送成功,這里返回的也是false,然爾我們并不能知道是具體的哪 一條消息發(fā)送失敗。

異步消息確認

異步的消息確認是通過一個監(jiān)聽器來實現(xiàn)的,當消息發(fā)送后,會接著執(zhí)行下面的邏輯,可能在稍會的一段時間,監(jiān)聽器監(jiān)聽到了Broker的返回,再進行邏輯的處理。

public static void main(String[] args) throws InterruptedException, TimeoutException, IOException {
        ConnectionFactory cf = new ConnectionFactory();
        cf.setHost("host");
        cf.setPort(5672);
        cf.setUsername("賬號");
        cf.setPassword("密碼");
        try(Connection connection = cf.newConnection();
            Channel channel = connection.createChannel()){
            channel.confirmSelect();
            ConfirmListener confirmListener = new ConfirmListener() {
                @Override
                public void handleAck(long deliveryTag, boolean multiple) throws IOException {
                    System.out.println("發(fā)送成功:" + deliveryTag + " multiple:" + multiple);
                }
                @Override
                public void handleNack(long deliveryTag, boolean multiple) throws IOException {
                    System.out.println("發(fā)送失敗:" + deliveryTag);
                }
            };
            channel.addConfirmListener(confirmListener);
            Map<String,String> mes = new HashMap<>();
            mes.put("name","11111");
            String messageStr = objectMapper.writeValueAsString(mes);
            for(int i = 0;i < 100;i++){
                channel.basicPublish(
                        "exchange.drinks",
                        "drinks.juzi",
                        null,
                        messageStr.getBytes());
            }
            Thread.sleep(Integer.MAX_VALUE);
        }
    }

當成功的發(fā)送消息的時候會回調(diào)監(jiān)聽器中的handleAck方法,如果沒有發(fā)送成功會回調(diào)handleNack方法 在這個監(jiān)聽器里面有兩個參數(shù)一個deliveryTagmultiple:

  • deliveryTag:表示當前的Channel發(fā)送的第幾條消息
  • multiple:是否在確認多條消息

這個異步的雖然在聽覺上感覺比較厲害些,這里也不推薦使用,原因和上面的一樣,我們并不能具休的知道是哪一條消息沒有被確認發(fā)送。

綜上:這里更加推薦單條消息確認,具體選擇哪一種還是要用業(yè)務(wù)做出選擇

注:注意一點是當一條消息成功的發(fā)送到Broker,但是如果沒有正確的路由到隊列,那么這時borker也是會返回true,因為Broker確時接收到了消息只是RoutingKey不可達,所以這里也會返回true,并且直接將消息丟棄

消息的返回機制

這個消息返回機制的作用就是在當一個消息成功的發(fā)送,但是并沒有正確路由到隊列的時候所回調(diào)的。

這也彌補了上面確認消息是否發(fā)送成功但沒有路由到隊列所返回true的問題 在使用消息返回機制的時候在發(fā)送消息時需要將mandatory置成true。再添加對應(yīng)的監(jiān)聽器。

public static void main(String[] args) throws InterruptedException, TimeoutException, IOException {
        ConnectionFactory cf = new ConnectionFactory();
        cf.setHost("host");
        cf.setPort(5672);
        cf.setUsername("賬號");
        cf.setPassword("密碼");
        try(Connection connection = cf.newConnection();
            Channel channel = connection.createChannel()){
            channel.addReturnListener(new ReturnCallback() {
                @Override
                public void handle(Return returnMessage) {
                    System.out.println("replyCode:" + returnMessage.getReplyCode() + " replyText:" + returnMessage.getReplyText() + " routingKey:"
                    + returnMessage.getRoutingKey() + " exchange:" + returnMessage.getExchange() + " body:" + new String(returnMessage.getBody()));
                }
            });
            Map<String,String> mes = new HashMap<>();
            mes.put("name","11111");
            String messageStr = objectMapper.writeValueAsString(mes);
            channel.basicPublish(
                    "exchange.drinks",
                    "drinks.juzi1",
                    true,
                    null,
                    messageStr.getBytes());
            Thread.sleep(Integer.MAX_VALUE);
        }
    }

這里的addReturnListener方法有兩個重載:只不過是handle的參數(shù)不同,一個是參數(shù)都顯示在了參數(shù)列表內(nèi),一個是將參數(shù)封裝到了Return對象內(nèi)。當handle被回調(diào)的時候也可以獲取到相應(yīng)的參數(shù)比如:exchange routingkey body。

注:保證消息可靠性投遞的前提是服務(wù)的高可用,服務(wù)不高可用談其它的都是扯

以上就是詳解RabbitMq如何做到消息的可靠性投遞的詳細內(nèi)容,更多關(guān)于RabbitMq 消息可靠性投遞的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Spring4如何自定義@Value功能詳解

    Spring4如何自定義@Value功能詳解

    這篇文章主要給大家介紹了關(guān)于Spring4如何自定義@Value功能的相關(guān)資料,文中通過示例代碼介紹的非常詳細,對大家學(xué)習(xí)或者使用spring4具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧。
    2017-09-09
  • springMVC返回Http響應(yīng)的實現(xiàn)

    springMVC返回Http響應(yīng)的實現(xiàn)

    本文主要介紹了在Spring Boot中使用@Controller、@ResponseBody和@RestController注解進行HTTP響應(yīng)返回的方法,具有一定的參考價值,感興趣的可以了解一下
    2025-03-03
  • Java8中方便又實用的Map函數(shù)總結(jié)

    Java8中方便又實用的Map函數(shù)總結(jié)

    java8之后,常用的Map接口中添加了一些非常實用的函數(shù),可以大大簡化一些特定場景的代碼編寫,提升代碼可讀性,快跟隨小編一起來看看吧
    2022-11-11
  • 解析ConcurrentHashMap:成員屬性、內(nèi)部類、構(gòu)造方法

    解析ConcurrentHashMap:成員屬性、內(nèi)部類、構(gòu)造方法

    ConcurrentHashMap是由Segment數(shù)組結(jié)構(gòu)和HashEntry數(shù)組結(jié)構(gòu)組成。Segment的結(jié)構(gòu)和HashMap類似,是一種數(shù)組和鏈表結(jié)構(gòu),今天給大家普及java面試常見問題---ConcurrentHashMap知識,一起看看吧
    2021-06-06
  • idea手動導(dǎo)入了包但編譯運行還是報找不到xxx.jar包的解決方案

    idea手動導(dǎo)入了包但編譯運行還是報找不到xxx.jar包的解決方案

    這篇文章主要介紹了idea手動導(dǎo)入了包但編譯運行還是報找不到xxx.jar包的解決方案,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2025-03-03
  • Springboot+ElementUi實現(xiàn)評論、回復(fù)、點贊功能

    Springboot+ElementUi實現(xiàn)評論、回復(fù)、點贊功能

    這篇文章主要介紹了通過Springboot ElementUi實現(xiàn)評論、回復(fù)、點贊功能。如果是自己評論的還可以刪除,刪除的規(guī)則是如果該評論下還有回復(fù),也一并刪除。需要的可以參考一下
    2022-01-01
  • 詳解eclipse項目中的.classpath文件原理

    詳解eclipse項目中的.classpath文件原理

    這篇文章介紹了eclipse項目中的.classpath文件的原理,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2021-12-12
  • 詳解Spring Cloud Gateway 數(shù)據(jù)庫存儲路由信息的擴展方案

    詳解Spring Cloud Gateway 數(shù)據(jù)庫存儲路由信息的擴展方案

    這篇文章主要介紹了詳解Spring Cloud Gateway 數(shù)據(jù)庫存儲路由信息的擴展方案,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-11-11
  • JVM中的flag設(shè)置詳解

    JVM中的flag設(shè)置詳解

    這篇文章主要介紹了JVM中的flag設(shè)置詳解,涉及堆大小設(shè)置,收集器設(shè)置等香公館內(nèi)容,小編覺得還是挺不錯的,具有一定借鑒價值,需要的朋友可以參考下
    2018-02-02
  • java 獲取一組數(shù)據(jù)中的最大值和最小值

    java 獲取一組數(shù)據(jù)中的最大值和最小值

    本文主要介紹了java 獲取一組數(shù)據(jù)中的最大值和最小值的方法。具有很好的參考價值,下面跟著小編一起來看下吧
    2017-02-02

最新評論

桦川县| 兴化市| 阿鲁科尔沁旗| 洞口县| 蒲江县| 揭阳市| 和平县| 黄冈市| 威海市| 新田县| 礼泉县| 海安县| 高要市| 盘锦市| 遵义县| 财经| 桂阳县| 皮山县| 密云县| 茂名市| 阿瓦提县| 康马县| 无极县| 双城市| 葫芦岛市| 张家港市| 荥经县| 罗田县| 凤阳县| 白山市| 青河县| 武鸣县| 河源市| 巩义市| 玉田县| 盐源县| 蒙阴县| 马边| 莎车县| 镇康县| 九龙县|