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

RabbitMQ保證消息的可靠性問題及解決

 更新時間:2026年04月29日 09:42:47   作者:又旅行又開拓的繩匠..  
文章主要討論了如何確保消息隊列中消息的可靠性,包括生產者重試機制、確認機制、MQ持久化配置和消費者確認機制等,還介紹了消費者本地重試和失敗消息入隊的處理策略,以及冪等性的實現(xiàn)方案,最后給出了一個綜合的配置示例

消息可靠性問題

在消息隊列,任何一個環(huán)節(jié)出問題都會導致消息的不可靠(消息丟失),如何確保消息的可靠性呢,需要考慮到其中的每個角色,生產者可靠性、MQ可靠性、消費者可靠性。

生產者可靠性

生產者重試

首先第一種情況,就是生產者發(fā)送消息時,出現(xiàn)了網絡故障,導致與MQ的連接中斷。

為了解決這個問題,SpringAMQP提供的消息發(fā)送時的重試機制。即:當RabbitTemplate與MQ連接超時后,多次重試。

spring:
  rabbitmq:
    connection-timeout: 1s # 設置MQ的連接超時時間
    template:
      retry:
        enabled: true # 開啟超時重試機制
        initial-interval: 1000ms # 失敗后的初始等待時間
        multiplier: 1 # 失敗后下次的等待時長倍數(shù),下次等待時長 = 上次等待時長 * multiplier
        max-attempts: 3 # 總共嘗試次數(shù)

耗盡重試次數(shù)后,依舊失敗,記錄失敗消息到數(shù)據庫失敗消息表,用于后期執(zhí)行補償錯誤。如使用定時任務去掃描這個表,重新發(fā)送消息

生產者確認

1.Publisher Return

消息投遞成功但路由失敗會調用Publisher Return回調方法返回異常信息。

2.Publisher Confirm

消息投遞成功返回ack,投遞失敗返回nack。

注意:消息投遞成功但可能路由失敗了,此時會通過Publisher Confirm返回ack,通過Publisher Return回調方法返回異常信息。

默認兩種機制都是關閉狀態(tài),需要通過配置文件來開啟。

spring:
  rabbitmq:
    publisher-confirm-type: correlated # 開啟publisher confirm機制,并設置confirm類型
    publisher-returns: true # 開啟publisher return機制

MQ可靠性

為了提升性能,默認情況下MQ的數(shù)據都是在內存存儲的臨時數(shù)據,重啟后就會消失。為了保證數(shù)據的可靠性,必須配置持久化,包括:

  • 交換機持久化
  • 隊列持久化
  • 消息持久化

在配置的時候默認都會持久化

消費者可靠性

消費者確認機制

為了確認消費者是否成功處理消息,RabbitMQ提供了消費者確認機制(Consumer Acknowledgement)。

即:當消費者處理消息結束后,應該向RabbitMQ發(fā)送一個回執(zhí),告知RabbitMQ消息處理狀態(tài)。

回執(zhí)有三種可選值:

  • ack:成功處理消息,RabbitMQ從隊列中刪除該消息
  • nack:消息處理失敗,RabbitMQ需要再次投遞消息
  • reject:消息處理失敗并拒絕該消息,RabbitMQ從隊列中刪除該消息

SpringAMQP幫我們實現(xiàn)了消息確認,并可以通過配置文件設置消息確認的處理方式,有三種模式:

none:不處理。即消息投遞給消費者后消息會立刻從MQ刪除。非常不安全,不建議使用

manual:手動模式。需要自己在業(yè)務代碼中調用api,發(fā)送ackreject,存在業(yè)務入侵,但更靈活

auto:自動模式。當業(yè)務正常執(zhí)行時則自動返回ack. 當業(yè)務出現(xiàn)異常時,根據異常判斷返回不同結果:

  • 如果是業(yè)務異常,會自動返回nack
  • 如果是消息處理或校驗異常,自動返回reject,返回的異常包括:MessageConversionException、MethodArgumentTypeMismatchException等

通過下面的配置可以修改消息確認的處理方式為auto:

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: auto # 自動ack

auto模式就是平常的寫法

manual模式需要手寫

  • Message:是spring AMQP封裝的底層消息對象。
  • Channel:是消費端與MQ基于通道的操作對象。
    @RabbitListener(queues = "simple.queue")
    public void listenSimpleQueueMessage(String msg, Channel channel, Message message) throws InterruptedException, IOException {
        log.info("spring 消費者接收到消息:【" + msg + "】");
        //返回nack
        //每個參數(shù)的意義:1.消息的標記 2.是否確認之前所有未確認的消息 3.是否重新入隊
        channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
//        log.info("消息處理完成");
//        //返回ack,每個參數(shù)的意義:1.消息的標記 2.是否確認之前所有消息
//        channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
    }

失敗重試機制

本地重試

當消費者出現(xiàn)異常后,消息會不斷requeue(重入隊)到隊列,再重新發(fā)送給消費者。如果消費者再次執(zhí)行依然出錯,消息會再次返回到隊列,再次投遞,直到消息處理成功為止。

極端情況就是消費者一直無法執(zhí)行成功,那么消息投遞就會無限循環(huán),導致mq的消息處理飆升,帶來不必要的壓力。

為了應對上述情況Spring又提供了消費者失敗重試機制:在消費者出現(xiàn)異常時利用本地重試,而不是無限制的投遞到mq隊列。

spring:
  rabbitmq:
    listener:
      simple:
        retry:
          enabled: true # 開啟消費者失敗重試
          initial-interval: 1000ms # 初識的失敗等待時長為1秒
          multiplier: 1 # 失敗的等待時長倍數(shù),下次等待時長 = 上次等待時長 * multiplier
          max-attempts: 3 # 最大重試次數(shù)
  • 開啟本地重試時,消息處理過程中拋出異常,不會請求到隊列,而是在消費者本地重試
  • 重試達到最大次數(shù)后,Spring會返回reject,消息會被丟棄

失敗消息入隊

本地測試達到最大重試次數(shù)后,消息會被丟棄。這在某些對于消息可靠性要求較高的業(yè)務場景下,顯然不太合適了。

因此Spring允許我們自定義重試次數(shù)耗盡后的消息處理策略,這個策略是由MessageRecovery接口來定義的,它有3個不同實現(xiàn):

  • RejectAndDontRequeueRecoverer:重試耗盡后,直接reject,丟棄消息。默認就是這種方式
  • ImmediateRequeueMessageRecoverer:重試耗盡后,返回nack,消息重新入隊
  • RepublishMessageRecoverer:重試耗盡后,將失敗消息投遞到指定的交換機

比較優(yōu)雅的一種處理方案是RepublishMessageRecoverer,失敗后將消息投遞到一個固定交換機,通過交換機將消息轉發(fā)到失敗消息隊列,程序監(jiān)聽失敗消息隊列,接收到失敗消息,將失敗消息存入失敗消息表,通過定時任務進行處理。

//在consumer服務中定義處理失敗消息的交換機和隊列
@Bean
public DirectExchange errorMessageExchange(){
    return new DirectExchange("error.direct");
}
@Bean
public Queue errorQueue(){
    return new Queue("error.queue", true);
}
@Bean
public Binding errorBinding(Queue errorQueue, DirectExchange errorMessageExchange){
    return BindingBuilder.bind(errorQueue).to(errorMessageExchange).with("error");
}
//定義一個RepublishMessageRecoverer,指定失敗消息投遞交換機的名稱及routingkey
@Bean
public MessageRecoverer republishMessageRecoverer(RabbitTemplate rabbitTemplate){
    return new RepublishMessageRecoverer(rabbitTemplate, "error.direct", "error");
}

監(jiān)聽失敗消息隊列將失敗消息寫入數(shù)據庫中,由人工定期處理

業(yè)務冪等性

冪等性:在程序開發(fā)中,是指同一個業(yè)務,執(zhí)行一次或多次對業(yè)務狀態(tài)的影響是一致的。

在程序開發(fā)中,是指同一個業(yè)務,執(zhí)行一次或多次對業(yè)務狀態(tài)的影響是一致的。

例如:

  • 根據id刪除數(shù)據
  • 查詢數(shù)據

但數(shù)據的更新往往不是冪等的,如果重復執(zhí)行可能造成不一樣的后果。比如:

  • 取消訂單,恢復庫存的業(yè)務。如果多次恢復就會出現(xiàn)庫存重復增加的情況
  • 退款業(yè)務。重復退款對商家而言會有經濟損失。

所以,我們要盡可能避免業(yè)務被重復執(zhí)行,然而在實際業(yè)務場景中,由于意外經常會出現(xiàn)業(yè)務被重復執(zhí)行的情況。

例如:

  • 頁面卡頓時頻繁刷新導致表單重復提交
  • 服務間調用的重試
  • MQ消息的重復投遞

因此,我們必須想辦法保證消息處理的冪等性。

這里給出兩種方案:

  • 唯一消息ID
  • 業(yè)務狀態(tài)判斷

唯一消息ID思路非常簡單:

  • 每一條消息都生成一個唯一的id,與消息一起投遞給消費者。
  • 消費者接收到消息后處理自己的業(yè)務,業(yè)務處理成功后將消息ID保存到數(shù)據庫或Redis
  • 如果下次又收到相同消息,去數(shù)據庫或Redis查詢判斷是否存在,存在則為重復消息放棄處理。

業(yè)務判斷就是基于業(yè)務本身的邏輯或狀態(tài)來判斷是否是重復的請求,不同的業(yè)務場景判斷的思路也不一樣。

例如在支付通知案例中,處理消息的業(yè)務邏輯是把訂單狀態(tài)從未支付修改為已支付。因此我們就可以在執(zhí)行更新時判斷訂單狀態(tài)是否是未支付,如果不是則證明訂單已經被處理過,無需重復處理。

相比較而言,使用唯一消息ID的方案需要操作數(shù)據庫或Redis保存消息ID,所以更推薦使用業(yè)務判斷的方案。

1.創(chuàng)建交換機,隊列,消息進行持久化

2.生產者:

  • 開啟消息發(fā)送失敗的重試策略,設置重試次數(shù)和間隔比例,耗盡重試次數(shù)后,依舊失敗,記錄失敗消息到數(shù)據庫失敗消息表,用于后期執(zhí)行錯誤補償.如使用定時任務去掃描這個表,重新發(fā)送消息
  • 開啟confirm機制,保證消息正確到達交換機,到達返回ack,沒有到達返回nack,寫入數(shù)據庫,后期重試
  • 開啟return機制,保證消息正確到達隊列,沒有到達調用ReturnCallback,寫入數(shù)據庫,后期重試

3.消費者:

  • 開啟手動ack,讓消費者消費成功后,手動提交.使用Redis來記錄消費失敗的次數(shù),如果到達閾值,則記錄到數(shù)據庫,后期使用人工干預
  • 自動ack + 重試耗盡的失敗策略,定義錯誤交換機隊列,后期通過人工進行干預

總結

以上為個人經驗,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關文章

  • 結合Mybatis聊聊對SQL注入的見解

    結合Mybatis聊聊對SQL注入的見解

    這篇文章主要介紹了結合Mybatis聊聊對SQL注入的見解,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • java 代碼塊與靜態(tài)代碼塊加載順序

    java 代碼塊與靜態(tài)代碼塊加載順序

    這篇文章主要介紹了java 代碼塊與靜態(tài)代碼塊加載順序的相關資料,需要的朋友可以參考下
    2017-07-07
  • springboot多文件或者文件夾壓縮成zip的方法

    springboot多文件或者文件夾壓縮成zip的方法

    最近碰到個需要下載zip壓縮包的需求,于是我在網上找了下別人寫好的zip工具類,下面通過本文給大家分享springboot多文件或者文件夾壓縮成zip的方法,感興趣的朋友一起看看吧
    2024-07-07
  • java Gui實現(xiàn)肯德基點餐收銀系統(tǒng)

    java Gui實現(xiàn)肯德基點餐收銀系統(tǒng)

    這篇文章主要為大家詳細介紹了java Gui實現(xiàn)肯德基點餐收銀系統(tǒng),文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2019-01-01
  • Java使用oshi獲取當前服務器狀態(tài)cpu、內存、存儲等核心信息方式

    Java使用oshi獲取當前服務器狀態(tài)cpu、內存、存儲等核心信息方式

    OSHI是基于JNA的跨平臺系統(tǒng)硬件信息庫,無需額外依賴,支持Windows/Linux/macOS/UNIX等系統(tǒng),提供CPU、內存、磁盤、網絡、電池等實時監(jiān)控數(shù)據,適用于資源監(jiān)控及可視化,含GUI示例與相關項目鏈接
    2025-09-09
  • RedisTemplate和Redisson的區(qū)別

    RedisTemplate和Redisson的區(qū)別

    本文主要介紹了RedisTemplate和Redisson的區(qū)別,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2025-11-11
  • Java使用Condition控制線程通信的方法實例詳解

    Java使用Condition控制線程通信的方法實例詳解

    這篇文章主要介紹了Java使用Condition控制線程通信的方法,結合實例形式分析了使用Condition類同步檢測控制線程通信的相關操作技巧,需要的朋友可以參考下
    2019-09-09
  • SpringBoot項目nohup啟動運行日志過大的解決方案

    SpringBoot項目nohup啟動運行日志過大的解決方案

    這篇文章主要介紹了SpringBoot項目nohup啟動運行日志過大的解決方案,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-05-05
  • JAVA 格式化日期、時間的方法

    JAVA 格式化日期、時間的方法

    這篇文章主要介紹了JAVA 格式化日期、時間的方法,文中講解非常細致,代碼幫助大家更好的理解和學習,感興趣的朋友可以了解下
    2020-06-06
  • Spring Security添加二次認證的項目實踐

    Spring Security添加二次認證的項目實踐

    在用戶自動登錄后,可以通過對密碼進行二次校驗進而確保用戶的真實性,本文就來介紹一下Spring Security添加二次認證的項目實踐,具有一定的參考價值,感興趣的可以了解一下
    2023-12-12

最新評論

尚义县| 松溪县| 宁远县| 怀远县| 西吉县| 五大连池市| 右玉县| 囊谦县| 阜宁县| 安仁县| 武宁县| 大连市| 巨鹿县| 正阳县| 肇源县| 佛坪县| 凤山市| 赫章县| 航空| 沙洋县| 华安县| 东方市| 英吉沙县| 尼玛县| 大厂| 化州市| 五原县| 天长市| 罗山县| 定陶县| 民丰县| 宜黄县| 瑞昌市| 岳西县| 乌拉特中旗| 济源市| 林州市| 三明市| 泰顺县| 龙江县| 北流市|