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

Pulsar源碼徹底解決重復(fù)消費(fèi)問(wèn)題

 更新時(shí)間:2023年05月29日 11:29:53   作者:crossoverJie  
這篇文章主要為大家介紹了Pulsar源碼徹底解決重復(fù)消費(fèi)問(wèn)題,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

背景

最近真是和 Pulsar 杠上了,業(yè)務(wù)團(tuán)隊(duì)反饋說(shuō)是線上有個(gè)應(yīng)用消息重復(fù)消費(fèi)。

而且在測(cè)試環(huán)境是可以穩(wěn)定復(fù)現(xiàn)的,根據(jù)經(jīng)驗(yàn)來(lái)看一般能穩(wěn)定復(fù)現(xiàn)的都比較好解決。

定位問(wèn)題

接著便是定位問(wèn)題了,根據(jù)之前的經(jīng)驗(yàn)讓業(yè)務(wù)按照這幾種情況先排查一下:

通過(guò)排查:1,2可以排除了。

  • 沒(méi)有相關(guān)日志
  • 存在異常,但最外層也捕獲了,所以不管有無(wú)異常都會(huì) ACK。

第三個(gè)也在消費(fèi)的入口和提交消息出計(jì)算了時(shí)間,最終發(fā)現(xiàn)都是在2s左右 ACK 的。

偽代碼如下:

Consumer consumer = client.newConsumer()
                .subscriptionType(SubscriptionType.Shared)
                .enableRetry(true)
                .topic(topic)
                .ackTimeout(30, TimeUnit.SECONDS)
                .subscriptionName("my-sub")
                .messageListener(new MessageListener<byte[]>() {
                    @SneakyThrows
                    @Override
                    public void received(Consumer<byte[]> consumer, Message<byte[]> msg) {
                        log.info("msg_id{}",msg.getMessageId().toString());
                        TimeUnit.SECONDS.sleep(2);
                        consumer.acknowledge(msg);
                    }
                })
                .subscribe();

那這就很奇怪了,因?yàn)榇a里配置的 ackTimeout 是 30s,理論上來(lái)說(shuō)是不會(huì)存在超時(shí)導(dǎo)致消息重發(fā)的。

為了排除是否是超時(shí)引起的,直接將業(yè)務(wù)代碼注釋掉了,等于是消息收到后立即就 ACK,經(jīng)過(guò)測(cè)試發(fā)現(xiàn)這樣確實(shí)就沒(méi)有重復(fù)消費(fèi)了。

為了再次確認(rèn)是不是和 ackTimeout 有關(guān),直接將 .ackTimeout(30, TimeUnit.SECONDS) 注釋掉后測(cè)試,發(fā)現(xiàn)也沒(méi)有重復(fù)消費(fèi)了。

確認(rèn)原因

既然如此那一定是和這個(gè)配置有關(guān)了,但看代碼確實(shí)沒(méi)有超時(shí),為了定位具體原因只有去看 client 的源碼了。

這里簡(jiǎn)單梳理下消息的消費(fèi)的流程:

  • 根據(jù) .receiverQueueSize(1000) 的配置,默認(rèn)情況下 broker 會(huì)直接給客戶端推送 1000 條消息。
  • 客戶端將這 1000 條消息保存到內(nèi)部隊(duì)列中。
  • 如果使用同步消費(fèi) receive() 時(shí),本質(zhì)上就是去 take 這個(gè)內(nèi)部隊(duì)列。
  • 如果是使用的是 messageListener 異步消費(fèi)并配置 ackTimeout,每當(dāng)從隊(duì)列里獲得一條消息后便會(huì)把這條消息加入 UnAckedMessageTracker 內(nèi)部的一個(gè)時(shí)間輪中,定時(shí)檢測(cè)頂部是否存在消息,如果存在則會(huì)觸發(fā)重新投遞。
    4.1 加入時(shí)間輪后,異步調(diào)用我們自定義的事件,這個(gè)異步操作是提交到一個(gè)無(wú)界隊(duì)列中由單個(gè)線程依次排隊(duì)執(zhí)行(這點(diǎn)是這次問(wèn)題的關(guān)鍵)
  • 業(yè)務(wù) ACK 的時(shí)候會(huì)從時(shí)間輪中刪除消息,所以如果消息 ACK 的足夠快,在第四步就不會(huì)獲取到消息進(jìn)行重新投遞。

整體流程如上圖,代碼細(xì)節(jié)如下圖:

所以問(wèn)題的根本原因就是寫(xiě)入時(shí)間輪(UnAckedMessageTracker)開(kāi)始倒計(jì)時(shí)的線程和回調(diào)業(yè)務(wù)邏輯的不是同一個(gè)線程。

如果業(yè)務(wù)執(zhí)行耗時(shí),等到消息從那個(gè)單線程的無(wú)界隊(duì)列中取出來(lái)的時(shí)候很有可能已經(jīng)過(guò)了 ackTimeou 的時(shí)間,從而導(dǎo)致了超時(shí)重發(fā)。

也就是用戶所理解的 ackTimeout 周期(應(yīng)該進(jìn)入回調(diào)時(shí)候開(kāi)始計(jì)時(shí))和 SDK 實(shí)現(xiàn)的不一致造成的。

之后我再次確認(rèn)同樣的代碼換為同步消費(fèi)是沒(méi)有問(wèn)題的,不會(huì)導(dǎo)致重復(fù)消費(fèi):

while (true) {
Message msg = consumer.receive();
            log.info(
                    "consumer Message received: " + new String(msg.getData()) + msg.getMessageId().toString());
            TimeUnit.SECONDS.sleep(2);
            consumer.acknowledge(msg);    
}

查看代碼后發(fā)現(xiàn)同步代碼的獲取消息和加入 UnAckedMessageTracker 時(shí)間輪是同步的,也就不會(huì)出現(xiàn)超時(shí)的問(wèn)題。

總結(jié)

所以其實(shí) 是messageListener 異步消費(fèi)的 ackTimeout 的語(yǔ)義是有問(wèn)題的,需要將加入 UnAckedMessageTracker 處移動(dòng)到回調(diào)函數(shù)中同步調(diào)用。

我查看了最新的 2.11.x 版本的代碼依然沒(méi)有修復(fù),正準(zhǔn)備提個(gè) PR 切換到 master 時(shí)才發(fā)現(xiàn)已經(jīng)有相關(guān)的 PR 了,只是還沒(méi)有發(fā)版。

修復(fù)的背景和思路也是類(lèi)似的,具體參考:

https://github.com/apache/pul...

其實(shí)業(yè)務(wù)中并不推薦使用 ackTimeout 這個(gè)配置了,不好預(yù)估時(shí)間從而導(dǎo)致超時(shí),而且我相信大部分業(yè)務(wù)配置好 ackTImeout 后直到后續(xù)出問(wèn)題的時(shí)候才想起來(lái)要改。

所以干脆一開(kāi)始就不要使用。

在 go 版本的 SDK 中直接廢棄掉了這個(gè)參數(shù),推薦使用 nack API 替換。

以上就是Pulsar源碼徹底解決重復(fù)消費(fèi)問(wèn)題的詳細(xì)內(nèi)容,更多關(guān)于Pulsar重復(fù)消費(fèi)解決的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Java深入講解Bean作用域與生命周期

    Java深入講解Bean作用域與生命周期

    這篇文章主要介紹了淺談Spring中Bean的作用域和生命周期,從創(chuàng)建到消亡的完整過(guò)程,例如人從出生到死亡的整個(gè)過(guò)程就是一個(gè)生命周期。本文將通過(guò)示例為大家詳細(xì)講講,感興趣的可以學(xué)習(xí)一下
    2022-06-06
  • Spring之詳解bean的實(shí)例化

    Spring之詳解bean的實(shí)例化

    這篇文章主要介紹了Spring之詳解bean的實(shí)例化,文章內(nèi)容詳細(xì),簡(jiǎn)單易懂,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2023-01-01
  • java 使用Scanner類(lèi)接收從控制臺(tái)輸入的數(shù)據(jù)方式

    java 使用Scanner類(lèi)接收從控制臺(tái)輸入的數(shù)據(jù)方式

    這篇文章主要介紹了java 使用Scanner類(lèi)接收從控制臺(tái)輸入的數(shù)據(jù)方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-08-08
  • JAVA實(shí)現(xiàn)301永久重定向方法

    JAVA實(shí)現(xiàn)301永久重定向方法

    本篇文章給大家總結(jié)了JAVA中實(shí)現(xiàn)永久重定向的方法以及詳細(xì)代碼,對(duì)此有需要的朋友可以參考學(xué)習(xí)下。
    2018-04-04
  • java web手寫(xiě)實(shí)現(xiàn)分頁(yè)功能

    java web手寫(xiě)實(shí)現(xiàn)分頁(yè)功能

    這篇文章主要為大家詳細(xì)介紹了java web手寫(xiě)實(shí)現(xiàn)分頁(yè)功能的方法,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2020-02-02
  • Java讀取properties文件連接數(shù)據(jù)庫(kù)的方法示例

    Java讀取properties文件連接數(shù)據(jù)庫(kù)的方法示例

    這篇文章主要介紹了Java讀取properties文件連接數(shù)據(jù)庫(kù)的方法示例,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2019-04-04
  • Spring Boot集成ShedLock分布式定時(shí)任務(wù)的實(shí)現(xiàn)示例

    Spring Boot集成ShedLock分布式定時(shí)任務(wù)的實(shí)現(xiàn)示例

    ShedLock確保您計(jì)劃的任務(wù)最多同時(shí)執(zhí)行一次。如果一個(gè)任務(wù)正在一個(gè)節(jié)點(diǎn)上執(zhí)行,則它會(huì)獲得一個(gè)鎖,該鎖將阻止從另一個(gè)節(jié)點(diǎn)(或線程)執(zhí)行同一任務(wù)。
    2021-05-05
  • Java?Stream流以及常用方法操作實(shí)例

    Java?Stream流以及常用方法操作實(shí)例

    Stream是對(duì)Java中集合的一種增強(qiáng)方式,使用它可以將集合的處理過(guò)程變得更加簡(jiǎn)潔、高效和易讀,這篇文章主要介紹了Java?Stream流以及常用方法的相關(guān)資料,需要的朋友可以參考下
    2025-08-08
  • Java中二叉樹(shù)的先序、中序、后序遍歷以及代碼實(shí)現(xiàn)

    Java中二叉樹(shù)的先序、中序、后序遍歷以及代碼實(shí)現(xiàn)

    這篇文章主要介紹了Java中二叉樹(shù)的先序、中序、后序遍歷以及代碼實(shí)現(xiàn),一棵二叉樹(shù)是結(jié)點(diǎn)的一個(gè)有限集合,該集合或者為空,或者是由一個(gè)根節(jié)點(diǎn)加上兩棵別稱為左子樹(shù)和右子樹(shù)的二叉樹(shù)組成,需要的朋友可以參考下
    2023-11-11
  • java藍(lán)橋杯試題

    java藍(lán)橋杯試題

    這篇文章主要介紹了java藍(lán)橋杯試題,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2020-02-02

最新評(píng)論

乌鲁木齐县| 遂昌县| 航空| 永清县| 阜阳市| 灵台县| 行唐县| 平度市| 宁陵县| 闽侯县| 兴化市| 化德县| 皮山县| 裕民县| 贡觉县| 红原县| 政和县| 根河市| 桃园市| 高雄市| 丰都县| 大余县| 如东县| 赤峰市| 涿鹿县| 三穗县| 福建省| 保山市| 尤溪县| 阳城县| 宁南县| 庆安县| 三穗县| 泊头市| 建阳市| 贵州省| 柯坪县| 遵化市| 新巴尔虎左旗| 息烽县| 大埔区|