鎖隊(duì)列原理詳解:從環(huán)形緩沖區(qū)到序列號(hào)機(jī)制)
說(shuō)到 Java 并發(fā)編程里的隊(duì)列大部分人的第一反應(yīng)是LinkedBlockingQueue、ArrayBlockingQueue或者是ConcurrentLinkedQueue。但如果你做過(guò)真正的高性能后端服務(wù)或者深入研究過(guò) Java 面試中那些和并發(fā)相關(guān)的硬核問(wèn)題大概率會(huì)碰上一個(gè)名字Disruptor。我第一次接觸 Disruptor 還是在看 LMAX 架構(gòu)文章的時(shí)候當(dāng)時(shí)就被它“每個(gè)時(shí)鐘周期處理 600 萬(wàn)訂單”這種說(shuō)法震住了。后來(lái)自己在項(xiàng)目里用上它才明白這東西根本不是什么黑魔法它只是把“并發(fā)”這件事的底層邏輯換了一套設(shè)計(jì)思路。先直接說(shuō)結(jié)論Disruptor 是一個(gè)無(wú)鎖的、有界的、用于線程間數(shù)據(jù)傳遞的環(huán)形隊(duì)列實(shí)現(xiàn)。它不是 JDK 自帶的也不是基于鎖或者 CAS 循環(huán)重試的傳統(tǒng)隊(duì)列而是通過(guò)一系列非常樸素但極其嚴(yán)謹(jǐn)?shù)膬?nèi)存布局、消費(fèi)依賴和序列號(hào)管理機(jī)制把并發(fā)競(jìng)爭(zhēng)降到最低。這篇文章我就用 Java 開(kāi)發(fā)者的視角把 Disruptor 的原理拆開(kāi)講清楚。會(huì)涉及到它比 BlockingQueue 快在哪里、為什么是環(huán)形的、Sequence 和 SequenceBarrier 是干什么的、偽共享是怎么回事、以及實(shí)際使用中哪些坑是我自己踩過(guò)并且覺(jué)得必須提醒你的。如果你正準(zhǔn)備 Java 面試或者正在為高吞吐場(chǎng)景選型又或者只是單純想搞明白“無(wú)鎖隊(duì)列到底是怎么做到無(wú)鎖的”這篇都適合你。我不會(huì)堆砌源碼但會(huì)把每個(gè)核心機(jī)制用大白話加實(shí)操經(jīng)驗(yàn)講透。1. 傳統(tǒng)隊(duì)列的性能瓶頸到底在哪從鎖到偽共享的層層損耗在聊 Disruptor 之前必須先回答一個(gè)問(wèn)題我們平時(shí)用的LinkedBlockingQueue和ArrayBlockingQueue究竟慢在哪里很多人以為慢在 CAS 自旋上其實(shí)更隱蔽的瓶頸在鎖競(jìng)爭(zhēng)、內(nèi)存屏障和緩存行沖突這三件事上。1.1 鎖競(jìng)爭(zhēng)線程之間最昂貴的協(xié)商成本ArrayBlockingQueue的生產(chǎn)者和消費(fèi)者共用一把鎖ReentrantLock。當(dāng)一個(gè)生產(chǎn)者線程正在往隊(duì)列里放數(shù)據(jù)消費(fèi)者線程想取數(shù)據(jù)就必須等待鎖釋放。這個(gè)等待過(guò)程不只是“等著”那么簡(jiǎn)單還涉及線程的上下文切換、操作系統(tǒng)的調(diào)度、鎖的爭(zhēng)用。我舉個(gè)例子你感受一下假設(shè)生產(chǎn)者線程 T1 持鎖寫(xiě)入消費(fèi)者線程 T2 在鎖上被阻塞。T2 被喚醒后需要重新判斷隊(duì)列狀態(tài)這個(gè)喚醒和切換的過(guò)程在低并發(fā)時(shí)無(wú)所謂但一旦線程數(shù)超過(guò) CPU 核心數(shù)或者生產(chǎn)者消費(fèi)者交替非常頻繁鎖的開(kāi)銷就變成平方級(jí)增長(zhǎng)。LinkedBlockingQueue雖然用了兩把鎖takeLock 和 putLock但依然存在鎖競(jìng)爭(zhēng)。而且鏈表結(jié)構(gòu)還有一個(gè)致命問(wèn)題每個(gè)節(jié)點(diǎn)都是一個(gè)對(duì)象創(chuàng)建和銷毀節(jié)點(diǎn)都會(huì)產(chǎn)生 GC 壓力節(jié)點(diǎn)之間的內(nèi)存地址不連續(xù)CPU 緩存命中率也低。1.2 偽共享一個(gè)看似無(wú)關(guān)卻致命的性能殺手這是很多 Java 開(kāi)發(fā)者容易忽略的概念。CPU 緩存是以緩存行Cache Line為單位的通常一個(gè)緩存行是 64 字節(jié)。當(dāng)兩個(gè)線程修改的是不同變量但這兩個(gè)變量恰好落在同一個(gè)緩存行里CPU 就會(huì)強(qiáng)制這個(gè)緩存行在兩個(gè)核心之間反復(fù)同步造成不必要的性能損耗。這種行為就叫偽共享False Sharing。傳統(tǒng)隊(duì)列中隊(duì)列的頭尾指針、狀態(tài)字段往往挨在一起存放生產(chǎn)者修改尾指針時(shí)消費(fèi)者讀取頭指針?biāo)诘木彺嫘袝?huì)失效反過(guò)來(lái)也一樣。這種互相拖后腿的現(xiàn)象在高并發(fā)下會(huì)放大得非常明顯。1.3 傳統(tǒng)隊(duì)列的吞吐量實(shí)感我之前用LinkedBlockingQueue做過(guò)一個(gè)壓測(cè)單生產(chǎn)者單消費(fèi)者模式每條消息 100 字節(jié)大概跑到每秒三十萬(wàn)到五十萬(wàn)條就很難上去了。而同樣的環(huán)境用 Disruptor可以輕松突破每秒百萬(wàn)級(jí)。差距不是一星半點(diǎn)而是量級(jí)上的碾壓。所以 Disruptor 解決的并不是“數(shù)據(jù)結(jié)構(gòu)”層面的問(wèn)題而是從 CPU 緩存、內(nèi)存布局、線程協(xié)作模型這些更底層的東西入手重新設(shè)計(jì)了一套方案。理解了這一點(diǎn)你再去看 Disruptor 的各個(gè)機(jī)制思路就會(huì)非常清晰。2. 環(huán)形緩沖區(qū)為什么有界環(huán)形結(jié)構(gòu)比鏈表更適合高并發(fā)Disruptor 內(nèi)部的核心存儲(chǔ)結(jié)構(gòu)就是一個(gè)預(yù)分配的有界環(huán)形數(shù)組這個(gè)數(shù)組被稱為 RingBuffer。為什么偏偏是環(huán)形我當(dāng)時(shí)的理解是環(huán)形結(jié)構(gòu)天然支持內(nèi)存預(yù)分配和復(fù)用這對(duì)高性能場(chǎng)景來(lái)說(shuō)是決定性的優(yōu)勢(shì)。2.1 數(shù)組預(yù)分配徹底消除 GC 壓力和內(nèi)存碎片RingBuffer 在初始化的時(shí)候就會(huì)把整個(gè)數(shù)組的對(duì)象一次性創(chuàng)建好后續(xù)生產(chǎn)者發(fā)布數(shù)據(jù)時(shí)只需要把數(shù)據(jù)從外部拷貝進(jìn)預(yù)先分配好的槽位即可。這和鏈表隊(duì)列每次 new 一個(gè) Node 不同Disruptor 在整個(gè)生命周期中幾乎不產(chǎn)生任何垃圾對(duì)象。GC 壓力小了STWStop The World自然就少延遲就更穩(wěn)定。這一點(diǎn)在交易系統(tǒng)、游戲服務(wù)器這類對(duì)延遲極其敏感的場(chǎng)景里是致命的優(yōu)勢(shì)。你可以把 RingBuffer 理解成一個(gè)循環(huán)利用的停車(chē)場(chǎng)每個(gè)車(chē)位都是固定的車(chē)到了就直接停進(jìn)空位不需要臨時(shí)搭車(chē)棚鏈表隊(duì)列則是每次來(lái)一輛車(chē)就得現(xiàn)搭一個(gè)棚開(kāi)走了再拆掉來(lái)回折騰成本高。2.2 為什么不直接用數(shù)組加鎖用數(shù)組并不新鮮ArrayBlockingQueue底層也是數(shù)組。問(wèn)題是它沒(méi)用環(huán)形結(jié)構(gòu)它每次讀寫(xiě)都需要計(jì)算數(shù)組邊界并且通過(guò)鎖來(lái)維持線程安全。而 Disruptor 的做法是利用“讀寫(xiě)下標(biāo)永遠(yuǎn)單調(diào)遞增”這個(gè)數(shù)學(xué)規(guī)律讓每個(gè)線程只需要維護(hù)自己關(guān)心的序列號(hào)完全不需要依賴鎖來(lái)協(xié)調(diào)邊界。換句話說(shuō)RingBuffer 不是用來(lái)“防止越界”的它是用來(lái)讓生產(chǎn)者和消費(fèi)者通過(guò)序列號(hào)各取所需而數(shù)組的環(huán)形特性只是為了復(fù)用內(nèi)存。真正決定誰(shuí)可以寫(xiě)入哪個(gè)槽位、誰(shuí)可以讀取哪個(gè)槽位的是下面要講的序列號(hào)機(jī)制。2.3 RingBuffer 的大小為什么必須是 2 的次冪這里有一個(gè)實(shí)際使用中經(jīng)常被忽略的細(xì)節(jié)RingBuffer 的容量必須是 2 的 N 次方默認(rèn)值是 16384也就是 2 的 14 次方。原因有兩個(gè)第一取模運(yùn)算position sequence (bufferSize - 1)可以直接用位運(yùn)算替代取模運(yùn)算的速度比%快很多。第二序列號(hào)回繞的邊界判斷更容易實(shí)現(xiàn)。只要保證容量是 2 的冪任何大于容量的序列號(hào)都能通過(guò)掩碼快速映射到具體槽位。我自己剛上手時(shí)習(xí)慣性傳了個(gè) 10000結(jié)果運(yùn)行直接報(bào)錯(cuò)看了源碼才發(fā)現(xiàn)int required 1; while (required bufferSize) required 1;這行邏輯它會(huì)把非 2 次冪的容量強(qiáng)制向上取整到最近的 2 的次冪。知道這個(gè)以后我配置容量時(shí)都會(huì)精確選擇 1024、4096、8192 這類值避免不必要的內(nèi)存開(kāi)銷。3. 序列號(hào)機(jī)制無(wú)鎖并發(fā)的核心契約如果說(shuō) RingBuffer 是 Disruptor 的骨架那么 Sequence序列號(hào)就是血液。Disruptor 無(wú)鎖的關(guān)鍵在于每個(gè)生產(chǎn)者和消費(fèi)者都維護(hù)一個(gè)自己的 Sequence多個(gè)線程之間通過(guò)對(duì)比這些 Sequence 的數(shù)值來(lái)決定能否讀寫(xiě)槽位而不是通過(guò)鎖去競(jìng)爭(zhēng)資源。3.1 Sequence 對(duì)象為什么要做緩存行填充先看源碼里的Sequence類你會(huì)發(fā)現(xiàn)它內(nèi)部維護(hù)了一個(gè)volatile long value。但光用 volatile 還不夠Disruptor 給這個(gè)value前后都塞了一大堆protected long p1, p2, p3...的占位字段硬生生把 64 字節(jié)的緩存行填滿了。為什么要這么干就是為了解決我前面提到的偽共享問(wèn)題。你想想生產(chǎn)者的寫(xiě)入序列號(hào)寫(xiě)進(jìn) value 時(shí)如果這個(gè) value 和消費(fèi)者的讀取序列號(hào)恰好落在同一個(gè)緩存行那每次消費(fèi)者讀取它自己的序列號(hào)時(shí)都會(huì)因?yàn)樯a(chǎn)者的寫(xiě)入導(dǎo)致緩存行失效然后去內(nèi)存里重新拉取性能大打折扣。Disruptor 的做法就是給每個(gè) Sequence 對(duì)象加上 padding確保一個(gè)緩存行里只會(huì)存在一個(gè)熱點(diǎn)的 value 字段。有一點(diǎn)要說(shuō)明這種填充手段在不同 JDK 版本上有區(qū)別。Java 8 之前大家常用Contended注解或者手動(dòng)補(bǔ)位Java 8 之后 JVM 提供了更優(yōu)雅的jdk.internal.vm.annotation.Contended注解但默認(rèn)只在 JDK 內(nèi)部類上生效我們自己業(yè)務(wù)類要用的話得加 JVM 參數(shù)-XX:-RestrictContended。而 Disruptor 為了兼容性和穩(wěn)定性選擇手動(dòng)補(bǔ)位的方式這個(gè)細(xì)節(jié)如果你在面試中提到會(huì)非常加分。3.2 生產(chǎn)者的發(fā)布流程cursor 和 gating sequence生產(chǎn)者在寫(xiě)入數(shù)據(jù)時(shí)需要申請(qǐng)一個(gè)寫(xiě)入位置。這個(gè)寫(xiě)入位置是基于一個(gè)全局的cursor當(dāng)前已發(fā)布的最大序列號(hào)來(lái)計(jì)算的。流程大致是這樣的生產(chǎn)者根據(jù)自己的生產(chǎn)者序號(hào)生成器ProducerSequencer申請(qǐng)下一個(gè)可用的序列號(hào)。這個(gè)申請(qǐng)過(guò)程需要檢查消費(fèi)者是否跟得上自己。具體來(lái)說(shuō)要拿自己的下一個(gè)序列號(hào)減去消費(fèi)者的最小序列號(hào)gating sequence看看差值是否已經(jīng)超過(guò)了 RingBuffer 容量。如果消費(fèi)者消費(fèi)太慢生產(chǎn)者就自旋等待直到消費(fèi)者那邊推進(jìn)了序列號(hào)騰出空間如果空間足夠生產(chǎn)者直接發(fā)布數(shù)據(jù)并發(fā)布事件通過(guò)Sequence的set方法更新 cursor 的值同時(shí)使用內(nèi)存屏障保證之前寫(xiě)入的數(shù)據(jù)對(duì)消費(fèi)者可見(jiàn)。這個(gè)機(jī)制在設(shè)計(jì)上非常像操作系統(tǒng)的生產(chǎn)者消費(fèi)者模型只不過(guò)把鎖替換成了“序列號(hào)比較”。當(dāng)然它也有等待策略后面我會(huì)講。3.3 消費(fèi)者的消費(fèi)流程SequenceBarrier 的協(xié)調(diào)作用消費(fèi)者側(cè)沒(méi)有直接用鎖而是通過(guò)SequenceBarrier序列屏障來(lái)協(xié)調(diào)。每個(gè)消費(fèi)者內(nèi)部都有一個(gè)Sequence表示自己消費(fèi)到了哪個(gè)位置。當(dāng)消費(fèi)者想要拿下一批數(shù)據(jù)時(shí)它會(huì)先讀取SequenceBarrier里緩存的 cursor 值。這個(gè)讀取不是簡(jiǎn)單的“讀變量”而是通過(guò)內(nèi)存屏障和SequenceBarrier的waitFor機(jī)制實(shí)現(xiàn)的。waitFor會(huì)返回當(dāng)前可消費(fèi)的最大序列號(hào)然后消費(fèi)者從這個(gè)序列號(hào)范圍內(nèi)批量獲取事件。這里有個(gè)設(shè)計(jì)精妙的地方多個(gè)消費(fèi)者可以依賴同一個(gè) SequenceBarrierDisruptor 會(huì)在背后維護(hù)一個(gè)gating sequence等于說(shuō)消費(fèi)者們看到的是一個(gè)“已經(jīng)被所有前置消費(fèi)者處理完的最遠(yuǎn)進(jìn)度”。換句話說(shuō)每個(gè)消費(fèi)者只保證自己處理的數(shù)據(jù)不會(huì)超過(guò)所有依賴方已經(jīng)處理完的位置。這樣就構(gòu)成了一個(gè)無(wú)鎖的依賴消費(fèi)鏈。4. 消費(fèi)依賴圖一旦你搞懂依賴模型Disruptor 就通了一半剛開(kāi)始用 Disruptor 時(shí)我最困惑的不是 API 怎么寫(xiě)而是它怎么處理復(fù)雜的業(yè)務(wù)流程。比如一個(gè)訂單數(shù)據(jù)進(jìn)來(lái)后需要先做風(fēng)控校驗(yàn)然后并行做積分累計(jì)和消息推送最后再做數(shù)據(jù)落庫(kù)。這種菱形依賴在 Disruptor 里是怎么表達(dá)的答案是消費(fèi)依賴圖Consumer Dependency Graph和SequenceBarrier的組合。4.1 單消費(fèi)者與多消費(fèi)者的消費(fèi)模式區(qū)別Disruptor 提供了兩種事件消費(fèi)模式EventHandler每個(gè)事件都會(huì)被所有注冊(cè)的消費(fèi)者都處理一遍屬于廣播模式。適合多個(gè)模塊都需要同一份數(shù)據(jù)的場(chǎng)景。WorkHandler每個(gè)事件只會(huì)被一個(gè)消費(fèi)者處理屬于競(jìng)爭(zhēng)模式。適合負(fù)載均衡分發(fā)的場(chǎng)景。選擇哪種模式取決于你的業(yè)務(wù)語(yǔ)義。比如日志收集場(chǎng)景一條日志來(lái)了既想寫(xiě)入本地又想上報(bào)監(jiān)控用EventHandler更合適如果只是想把這些日志分發(fā)到 Kafka那用WorkHandler更合適。它們底層的消費(fèi)者序列號(hào)管理邏輯不太一樣競(jìng)爭(zhēng)模式下 Disruptor 內(nèi)部會(huì)自動(dòng)為多個(gè) WorkProcessor 維護(hù)同一個(gè) WorkSequence確保事件不會(huì)重復(fù)分配。4.2 依賴鏈路的構(gòu)建SequenceBarrier 的層級(jí)關(guān)系在實(shí)際代碼中構(gòu)建依賴關(guān)系需要使用多個(gè)SequenceBarrier。簡(jiǎn)單來(lái)說(shuō)如果你想讓事件消費(fèi) A 必須發(fā)生在 B、C 并行處理之前那 B 和 C 的 SequenceBarrier 就會(huì)各自依賴 A 的 Sequence而如果有一個(gè) D 必須等 B 和 C 都完才處理那 D 的 SequenceBarrier 依賴的就是 B 和 C 的最小序列號(hào)即兩者中處理得最慢的那個(gè)位置。這里有一個(gè)容易犯迷糊的點(diǎn)Disruptor 的依賴是“多消費(fèi)者序列號(hào)的集合”而不是單個(gè)消費(fèi)者。所以在構(gòu)建BatchEventProcessor時(shí)每個(gè)消費(fèi)者都可以持有任意多個(gè)上游消費(fèi)者的 Sequence 作為門(mén)閂只有當(dāng)所有上游都推進(jìn)到某個(gè)位置下游才可以消費(fèi)對(duì)應(yīng)位置的事件。這種設(shè)計(jì)比用鎖或者 ConcurrentHashMap 做狀態(tài)同步要高效得多因?yàn)槿讨皇菙?shù)值比較沒(méi)有任何阻塞點(diǎn)。4.3 菱形依賴的代碼示意與邊界用代碼來(lái)看假設(shè)有三個(gè)消費(fèi)者EventHandlerOrderEvent riskCheck (event, sequence, endOfBatch) - doRiskCheck(event); EventHandlerOrderEvent pointsAccum (event, sequence, endOfBatch) - doAccumulate(event); EventHandlerOrderEvent pushNotify (event, sequence, endOfBatch) - doPush(event); EventHandlerOrderEvent saveDb (event, sequence, endOfBatch) - doSave(event);如果希望風(fēng)控校驗(yàn)完成之后再并行執(zhí)行積分累計(jì)和推送最后數(shù)據(jù)入庫(kù)構(gòu)建依賴時(shí)就要利用Disruptor的after方法EventHandlerGroupOrderEvent groupAfterRisk disruptor.after(riskCheck); groupAfterRisk.handleEventsWith(pointsAccum, pushNotify); groupAfterRisk.then(saveDb);注意then方法返回的是EventHandlerGroup并且它內(nèi)部會(huì)把pointsAccum和pushNotify的序列集合作為下游屏障。整體上的效果就是saveDb 永遠(yuǎn)不會(huì)越過(guò) pointsAccum 和 pushNotify 的最小進(jìn)度去消費(fèi)事件。如果你的業(yè)務(wù)在消費(fèi)依賴上遇到了“某個(gè)事件必須等兩個(gè)并行任務(wù)都完成才能繼續(xù)”的場(chǎng)景這個(gè)模型就是為你設(shè)計(jì)的。5. 發(fā)布流程中的三個(gè)關(guān)鍵步驟從事件轉(zhuǎn)換到最終發(fā)布真正動(dòng)手寫(xiě) Disruptor 生產(chǎn)者代碼你會(huì)發(fā)現(xiàn)發(fā)布流程其實(shí)就三步獲取槽位、寫(xiě)入數(shù)據(jù)、發(fā)布事件。但每一步背后都有值得展開(kāi)的機(jī)制和容易出錯(cuò)的細(xì)節(jié)。5.1 translate 階段利用 EventTranslator 干臟活累活Disruptor 推薦通過(guò)EventTranslator或者EventTranslatorOneArg來(lái)把業(yè)務(wù)數(shù)據(jù)寫(xiě)入 RingBuffer 的預(yù)分配槽位中。比如這樣EventTranslatorOneArgOrderEvent, Order TRANSLATOR (event, sequence, order) - { event.setId(order.getId()); event.setPrice(order.getPrice()); event.setTimestamp(order.getTimestamp()); }; ringBuffer.publishEvent(TRANSLATOR, order);publishEvent內(nèi)部會(huì)先申請(qǐng)序列號(hào)sequence然后調(diào)用translator.translateTo(event, sequence, order)再走發(fā)布流程。這個(gè)設(shè)計(jì)從使用者的角度來(lái)看很舒服你完全不用關(guān)心怎么拿序列號(hào)、怎么處理槽位競(jìng)爭(zhēng)只需要把業(yè)務(wù)數(shù)據(jù)映射到事件對(duì)象上即可。每個(gè)translateTo調(diào)用都會(huì)拿到一個(gè)對(duì)應(yīng)的槽位索引但如果你定義的事件對(duì)象是有狀態(tài)的比如可復(fù)用對(duì)象就必須注意把舊值清干凈否則會(huì)出現(xiàn)臟數(shù)據(jù)串?dāng)_。這是我踩過(guò)的一個(gè)很典型的坑事件對(duì)象內(nèi)有 list 字段第二次發(fā)布時(shí)忘了 clear導(dǎo)致消息內(nèi)容殘留。5.2 發(fā)布的內(nèi)存屏障保證其他線程一定能看到寫(xiě)入的數(shù)據(jù)發(fā)布事件時(shí)最關(guān)鍵的一步是ringBuffer.publish(sequence)。這一步會(huì)調(diào)用Sequencer的publish方法內(nèi)部重點(diǎn)在于對(duì)cursor的更新同時(shí)確保之前所有寫(xiě)入操作按順序?qū)οM(fèi)者可見(jiàn)。這個(gè)語(yǔ)義依賴的是 Java 的 volatile 變量寫(xiě)和讀之間的 happens-before 關(guān)系。我在實(shí)際項(xiàng)目中曾經(jīng)試圖使用普通變量來(lái)寫(xiě) RingBuffer 里的事件字段以為發(fā)布時(shí)不寫(xiě) volatile 也能靠后續(xù)的原子操作兜底結(jié)果消費(fèi)者端出現(xiàn)了偶發(fā)讀到空值的問(wèn)題。后來(lái)老老實(shí)實(shí)遵循 Disruptor 的寫(xiě)法所有數(shù)據(jù)先寫(xiě)進(jìn)預(yù)分配槽位再統(tǒng)一發(fā)布問(wèn)題消失。這種“先寫(xiě)數(shù)據(jù)、再發(fā)布”的順序非常關(guān)鍵Disruptor 管它叫做“Memory Barrier”。你只需要記住任何對(duì) RingBuffer 中事件字段的修改必須在調(diào)用publish之前完成不要反過(guò)來(lái)。5.3 多生產(chǎn)者場(chǎng)景下序列號(hào)的分配AtomicLong 與緩存行填充說(shuō)到多生產(chǎn)者就繞不開(kāi)MultiProducerSequencer。在多生產(chǎn)者模式下多個(gè)線程同時(shí)申請(qǐng)序列號(hào)Disruptor 內(nèi)部使用了一個(gè)AtomicLong通過(guò) CAS 自旋來(lái)管理cursor的分配。每次生產(chǎn)者申請(qǐng)序列號(hào)long current cursor.get(); long next current 1; while (!cursor.compareAndSet(current, next)) { current cursor.get(); next current 1; }這就是一個(gè)標(biāo)準(zhǔn) CAS 循環(huán)。這里看似還是存在競(jìng)爭(zhēng)但競(jìng)爭(zhēng)的粒度和鎖完全不同CAS 競(jìng)爭(zhēng)的是一個(gè) 8 字節(jié)的變量而且失敗后線程不會(huì)掛起只是自旋重試成本遠(yuǎn)低于鎖。再加上原子類內(nèi)部也做了緩存行填充多個(gè)生產(chǎn)者線程修改同一個(gè) AtomicLong 的性能表現(xiàn)遠(yuǎn)好于預(yù)期。單生產(chǎn)者模式下則完全不同它只需要一個(gè)普通變量加內(nèi)存屏障就可以安全發(fā)布因?yàn)楦緵](méi)有競(jìng)爭(zhēng)。所以選型時(shí)一定要誠(chéng)實(shí)評(píng)估自己的場(chǎng)景單生產(chǎn)者單消費(fèi)者、單生產(chǎn)者多消費(fèi)者、多生產(chǎn)者多消費(fèi)者分別對(duì)應(yīng)完全不同的內(nèi)部實(shí)現(xiàn)和生產(chǎn)效率。6. 等待策略的選擇無(wú)鎖不等于零等待關(guān)鍵看你愿意用 CPU 換什么很多人以為 Disruptor 無(wú)鎖那就意味著消費(fèi)者永遠(yuǎn)在忙等、CPU 消耗極高。實(shí)際上 Disruptor 提供了多種等待策略它們之間的區(qū)別本質(zhì)上是“CPU 資源”和“延遲”之間的權(quán)衡。搞不清這一點(diǎn)就亂選策略生產(chǎn)環(huán)境丟消費(fèi)速度和延遲指標(biāo)是遲早的事。6.1 四種常用等待策略對(duì)比我先列一個(gè)基于實(shí)際壓測(cè)經(jīng)驗(yàn)的表格方便你直觀對(duì)比。等待策略適用場(chǎng)景CPU 占用延遲表現(xiàn)我的建議BusySpinWaitStrategy消費(fèi)者線程數(shù)不超過(guò) CPU 核心數(shù)且線程長(zhǎng)期活躍高最低專用于超低延遲場(chǎng)景比如高頻交易YieldingWaitStrategy競(jìng)爭(zhēng)激烈但希望保留一部分 CPU 給其他任務(wù)中高低大部分高并發(fā)場(chǎng)景首選SleepingWaitStrategy對(duì)延遲不那么敏感但想省 CPU低中高適合日志異步批量上報(bào)BlockingWaitStrategy線程會(huì)被掛起適合對(duì) CPU 資源極度敏感最低最高謹(jǐn)慎使用延遲抖動(dòng)明顯單看這張表你可能還是會(huì)猶豫我以自己的經(jīng)驗(yàn)補(bǔ)充一點(diǎn)如果你的延遲要求是亞毫秒級(jí)別用BlockingWaitStrategy它內(nèi)部的鎖競(jìng)爭(zhēng)會(huì)直接毀掉 Disruptor 的架構(gòu)優(yōu)勢(shì)如果只是需要低 CPU 占用并且能接受幾毫秒延遲SleepingWaitStrategy是合理選擇。6.2 等待策略背后的小設(shè)計(jì)缺陷和注意事項(xiàng)一個(gè)容易出問(wèn)題的點(diǎn)是YieldingWaitStrategy。它內(nèi)部使用Thread.yield()讓出 CPU但yield其實(shí)不保證一定會(huì)讓出而且依賴 JVM 實(shí)現(xiàn)。在高負(fù)載下如果大量消費(fèi)者同時(shí)調(diào)用 yield線程調(diào)度的開(kāi)銷可能反而比自旋還大。我壓測(cè)時(shí)曾把消費(fèi)者數(shù)量設(shè)為 12機(jī)器只有 8 核結(jié)果整體吞吐反而下降后來(lái)改成SleepingWaitStrategy才穩(wěn)定下來(lái)。BusySpinWaitStrategy是性能最好的但前提是消費(fèi)者線程真正的“釘”在 CPU 上。假如消費(fèi)者線程偶爾會(huì)被其他業(yè)務(wù)代碼搶走自旋就變成無(wú)效空轉(zhuǎn)CPU 白燒。此時(shí)你會(huì)看到 CPU 飆高但沒(méi)有吞吐提升。所以選等待策略要跟線程綁定、核心數(shù)結(jié)合來(lái)看不要單看一個(gè)指標(biāo)。7. Disruptor 里的常見(jiàn)誤解無(wú)鎖、并行、性能幻覺(jué)每當(dāng)我跟同事聊 Disruptor 時(shí)都能聽(tīng)到各種想當(dāng)然的說(shuō)法。這里我把最典型的幾個(gè)誤解單獨(dú)拎出來(lái)用實(shí)際經(jīng)驗(yàn)說(shuō)明一下幫你也避開(kāi)這些坑。7.1 誤解一無(wú)鎖就是零阻塞、零等待完全不是。Disruptor 的無(wú)鎖是指不使用鎖作為并發(fā)協(xié)調(diào)手段但消費(fèi)者如果消費(fèi)速度跟不上生產(chǎn)者生產(chǎn)者會(huì)通過(guò)自旋等待或者等待策略被“限速”。這種自旋等待雖然不像鎖那樣讓線程休眠但仍然是一種阻塞。區(qū)別在于自旋等待不會(huì)導(dǎo)致線程上下文切換成本遠(yuǎn)低于鎖。也就是說(shuō)Disruptor 能扛住瞬時(shí)大量事件積壓但如果你一直讓生產(chǎn)者超速生產(chǎn)消費(fèi)者依然會(huì)形成背壓只是這種背壓更平滑、CPU 消耗更可控。7.2 誤解二EventHandler 越多消費(fèi)速度越快這是最常踩的坑。Disruptor 的EventHandler默認(rèn)是廣播模式多個(gè)處理器處理同一條數(shù)據(jù)的場(chǎng)景下每個(gè)消費(fèi)者都會(huì)拿到所有事件所以增加EventHandler并不會(huì)提升單條消息的處理吞吐而是增加處理鏈路的并行能力。如果想真正提速應(yīng)該把處理任務(wù)分片使用WorkHandler或者自己實(shí)現(xiàn)多個(gè)處理線程競(jìng)爭(zhēng)消費(fèi)。我見(jiàn)過(guò)一個(gè)新人把同樣邏輯的 EventHandler 重復(fù)注冊(cè)了三個(gè)以為能并發(fā)處理提升三倍速度結(jié)果所有事件被重復(fù)執(zhí)行了三次差點(diǎn)產(chǎn)生扣款重復(fù)。如果你也準(zhǔn)備用 WorkHandler務(wù)必記住它的消費(fèi)邏輯必須是冪等的否則重復(fù)消費(fèi)會(huì)變成大事故。7.3 誤解三Disruptor 應(yīng)該用來(lái)替代 Kafka這其實(shí)是完全不同的兩種東西。Disruptor 是進(jìn)程內(nèi)的內(nèi)存隊(duì)列數(shù)據(jù)不跨節(jié)點(diǎn)、不持久化進(jìn)程一崩數(shù)據(jù)全丟Kafka 是分布式消息中間件具備持久化、分區(qū)、副本、跨機(jī)容災(zāi)能力。它們解決的完全不是一個(gè)層面的問(wèn)題。Disruptor 的定位更像是 ConcurrentLinkedQueue 和 ArrayBlockingQueue 的高性能替代品是應(yīng)用內(nèi)部的管道。Kafka 這層屬于服務(wù)間通信。你完全可以也可以在業(yè)務(wù)里把兩者結(jié)合Disruptor 做應(yīng)用內(nèi)的異步削峰Kafka 做服務(wù)間的事件投遞。8. 實(shí)際工程中的選型建議與一套可落地的示例講了一堆原理最后還是回到工程落地。Disruptor 不是萬(wàn)金油它有自己的適用邊界。盲目的把系統(tǒng)里所有隊(duì)列都換成 Disruptor 是不理智的。我根據(jù)自己的項(xiàng)目經(jīng)驗(yàn)總結(jié)一套可復(fù)用的決策思路和一個(gè)完整的代碼骨架。8.1 什么場(chǎng)景適合上 Disruptor什么場(chǎng)景別用我的經(jīng)驗(yàn)是核心指標(biāo)是吞吐量和延遲抖動(dòng)且數(shù)據(jù)結(jié)構(gòu)相對(duì)固定、業(yè)務(wù)處理很快適合用 Disruptor。典型場(chǎng)景如訂單處理流水線、行情數(shù)據(jù)分發(fā)、日志異步批量寫(xiě)入。反過(guò)來(lái)如果你需要消息持久化、需要分布式消費(fèi)組、需要消息積壓觸達(dá)百萬(wàn)級(jí)那直接選 MQ 中間件別拿 Disruptor 硬扛。如果業(yè)務(wù)數(shù)據(jù)的到達(dá)模式極不均勻且消費(fèi)者處理速度波峰波谷巨大也要慎重因?yàn)?Disruptor 的預(yù)分配緩沖會(huì)一直占著內(nèi)存。另外還有一個(gè)很容易忽略的點(diǎn)Disruptor 適合“管道化處理”如果事件處理邏輯極其復(fù)雜且依賴大量不可控外部調(diào)用比如遠(yuǎn)程 HTTP那么消費(fèi)者線程很容易變成性能瓶頸。這不是 Disruptor 的問(wèn)題而是你的處理任務(wù)太重。真要上也讓消費(fèi)者內(nèi)部再用線程池去異步化別再同步阻塞。8.2 一個(gè)可以直接套用的單生產(chǎn)者多消費(fèi)者示例我把最核心的單生產(chǎn)者多消費(fèi)者示例寫(xiě)一下包含完整的初始化、發(fā)布、銷毀過(guò)程注釋會(huì)比較全方便你直接抄作業(yè)。public class OrderEvent { private long id; private double price; private long timestamp; // getters/setters 省略 } public class OrderEventFactory implements EventFactoryOrderEvent { Override public OrderEvent newInstance() { return new OrderEvent(); } } public class OrderEventHandler implements EventHandlerOrderEvent { private String consumerName; public OrderEventHandler(String consumerName) { this.consumerName consumerName; } Override public void onEvent(OrderEvent event, long sequence, boolean endOfBatch) { // 這里就是消費(fèi)者真正處理事件的入口 System.out.println(consumerName 消費(fèi)事件: id event.getId() , price event.getPrice() , seq sequence); } }啟動(dòng)以及發(fā)布的核心代碼如下// 1. 初始化 Disruptor int bufferSize 1024; DisruptorOrderEvent disruptor new Disruptor( new OrderEventFactory(), bufferSize, Executors.defaultThreadFactory(), ProducerType.SINGLE, new YieldingWaitStrategy() ); // 2. 注冊(cè)消費(fèi)者 disruptor.handleEventsWith( new OrderEventHandler(consumerA), new OrderEventHandler(consumerB) ); // 3. 啟動(dòng) disruptor.start(); // 4. 獲取 RingBuffer RingBufferOrderEvent ringBuffer disruptor.getRingBuffer(); // 5. 在業(yè)務(wù)線程中發(fā)布事件 EventTranslatorOneArgOrderEvent, Order translator (event, sequence, order) - { event.setId(order.getId()); event.setPrice(order.getPrice()); event.setTimestamp(System.currentTimeMillis()); }; for (Order order : orders) { ringBuffer.publishEvent(translator, order); }如果要用 WorkHandler 實(shí)現(xiàn)負(fù)載均衡只需要把 handleEventsWith 換成 handleEventsWithWorkerPooldisruptor.handleEventsWithWorkerPool( new OrderWorkHandler(consumerA), new OrderWorkHandler(consumerB) );關(guān)鍵是記住不同模式注冊(cè) API 不一樣語(yǔ)義也差很多代碼很容易跑通但邏輯可能不是你要的。8.3 消費(fèi)完成后的資源釋放與優(yōu)雅停機(jī)Disruptor 用完后需要優(yōu)雅關(guān)閉很多線上故障都出現(xiàn)在重啟和停機(jī)階段。標(biāo)準(zhǔn)做法是調(diào)用disruptor.shutdown()它會(huì)等待所有注冊(cè)的事件處理器處理完當(dāng)前 RingBuffer 中已發(fā)布的事件然后才返回。如果你設(shè)置了超時(shí)時(shí)間也可以用shutdown(long timeout, TimeUnit unit)。另一個(gè)容易被忽略的點(diǎn)是事件體本身是復(fù)用的所以在停機(jī)時(shí)把 RingBuffer 里剩余事件對(duì)象中的敏感數(shù)據(jù)清掉防止內(nèi)存中堆積臟數(shù)據(jù)。對(duì)安全要求高的場(chǎng)景比如交易訂單這一點(diǎn)特別重要?jiǎng)e嫌麻煩。8.4 監(jiān)控和性能調(diào)優(yōu)的落地建議Disruptor 部署到生產(chǎn)環(huán)境后不可能不監(jiān)控。我自己習(xí)慣重點(diǎn)觀察這幾個(gè)指標(biāo)RingBuffer 剩余容量如果長(zhǎng)期低于容量的 10%說(shuō)明消費(fèi)者處理不過(guò)來(lái)。每個(gè)消費(fèi)者 Sequence 與 cursor 的差值差值長(zhǎng)期大于容量的一半就說(shuō)明消費(fèi)滯后嚴(yán)重。事件處理耗時(shí)分布可以使用 Micrometer 這類工具記錄onEvent耗時(shí)觀察 P99 和 P99.9。壓測(cè)時(shí)建議用JMH寫(xiě)基準(zhǔn)測(cè)試把吞吐量和延遲一起看。只看吞吐量不看延遲是自欺欺人因?yàn)橛械牡却呗詾榱送掏驴梢誀奚艽蟮难舆t抖動(dòng)。9. 面試中的 Disruptor 考點(diǎn)串講如果你是為了準(zhǔn)備 Java 面試點(diǎn)進(jìn)來(lái)的這一節(jié)專門(mén)為你服務(wù)。Disruptor 在面試中算是一個(gè)比較進(jìn)階但不冷門(mén)的題懂的候選人通常會(huì)給面試官留下“底層扎實(shí)”的印象。常見(jiàn)的問(wèn)題有這些我附上最精煉的回答思路“Disruptor 為什么不需要鎖” 核心是用序列號(hào)加內(nèi)存屏障管理并發(fā)避免線程掛起和上下文切換?!癉isruptor 是如何解決偽共享的” 每個(gè) Sequence 做緩存行填充讓熱字段獨(dú)占緩存行?!癛ingBuffer 為什么比鏈表性能高” 數(shù)組內(nèi)存連續(xù)性更好、預(yù)分配對(duì)象無(wú) GC、索引計(jì)算可以用位運(yùn)算?!岸嗌a(chǎn)者和單生產(chǎn)者的區(qū)別” 多生產(chǎn)者需要 CAS 分配序列號(hào)單生產(chǎn)者只需要一個(gè)變量加內(nèi)存屏障?!癉isruptor 怎么實(shí)現(xiàn)依賴消費(fèi)” 通過(guò) SequenceBarrier 持有上游消費(fèi)者的 Sequence 集合取最小值做門(mén)檻。如果你能把這些機(jī)制用自己的語(yǔ)言講清再結(jié)合一次實(shí)際壓測(cè)數(shù)據(jù)面試官基本就很難在這一塊把你問(wèn)倒了。不過(guò)面試歸面試真正重要的是把原理理解透然后應(yīng)用到你的實(shí)際業(yè)務(wù)中。我最后的體會(huì)是Disruptor 最大的價(jià)值不僅在于“快”更在于它提供了一種和傳統(tǒng)并發(fā)思維完全不同的視角——通過(guò)設(shè)計(jì)避免競(jìng)爭(zhēng)而不是通過(guò)協(xié)調(diào)解決競(jìng)爭(zhēng)。項(xiàng)目里如果能找到合適的契合點(diǎn)它帶來(lái)的穩(wěn)定性和可預(yù)測(cè)延遲會(huì)讓后端系統(tǒng)的整體質(zhì)量上一個(gè)臺(tái)階。