者-消費(fèi)者模式與并行任務(wù)調(diào)度:從BlockingQueue到虛擬線程的工程實(shí)踐)
我想先把這次做的東西說清楚——這個(gè)項(xiàng)目圍繞的是生產(chǎn)者-消費(fèi)者模式、并行任務(wù)調(diào)度以及一個(gè)經(jīng)常被忽略的細(xì)節(jié)更簡(jiǎn)潔的注釋和每項(xiàng)改進(jìn)的詳細(xì)解釋。我自己維護(hù)過一套高吞吐的通知推送組件早期代碼就是“能跑就行”的水平隊(duì)列選型靠拍腦袋線程池參數(shù)全是默認(rèn)值注釋寫了等于沒寫后來踩了數(shù)據(jù)丟失、線程阻塞、代碼完全沒人敢改的坑才回過頭來用這套思路徹底重構(gòu)了一遍。適合讀這篇文章的人有兩類一類是剛接觸并發(fā)編程想知道 BlockingQueue、線程池、虛擬線程到底各自解決什么問題另一類是已經(jīng)寫了幾個(gè)月生產(chǎn)者-消費(fèi)者代碼但總覺得哪里別扭想看看“規(guī)范版本”長(zhǎng)什么樣。這篇內(nèi)容不會(huì)只丟給你一段能跑的代碼而是把我踩過的坑、每一步改進(jìn)背后的原因、包括注釋為什么要那樣寫全部拆開講清楚。1. 先搞清楚生產(chǎn)者-消費(fèi)者模式在工程里到底解決什么問題1.1 生產(chǎn)慢、消費(fèi)慢怎么讓他們別互相拖累很多人在面試?yán)锉尺^生產(chǎn)者-消費(fèi)者模式的定義但真到了業(yè)務(wù)里未必知道它到底在解什么題。我用一個(gè)日常場(chǎng)景說明你去咖啡店點(diǎn)單收銀員是生產(chǎn)者咖啡師是消費(fèi)者中間那個(gè)放訂單的臺(tái)子就是隊(duì)列。如果臺(tái)子太小收銀員每一次都要等咖啡師做完一杯才敢接下一單點(diǎn)單速度被拖死如果臺(tái)子太大訂單堆得到處都是消費(fèi)者又忙不過來。生產(chǎn)者-消費(fèi)者模式要做的事就是在“產(chǎn)生任務(wù)”和“執(zhí)行任務(wù)”中間加一層緩沖讓兩邊速度不匹配時(shí)不至于互相阻塞。工程里最常見的應(yīng)用場(chǎng)景是這幾類消息推送業(yè)務(wù)系統(tǒng)產(chǎn)生推送請(qǐng)求推送服務(wù)異步消費(fèi)高峰期削峰日志收集應(yīng)用線程寫日志到隊(duì)列后臺(tái)線程批量刷盤批量寫庫上游產(chǎn)生大量變更下游攢一批再寫減少數(shù)據(jù)庫連接壓力爬蟲任務(wù)解析抓取任務(wù)一條條進(jìn)隊(duì)列解析線程并行處理訂單狀態(tài)通知訂單系統(tǒng)產(chǎn)生事件通知服務(wù)消費(fèi)并發(fā)送短信/郵件你會(huì)發(fā)現(xiàn)它們的共同點(diǎn)生產(chǎn)端和消費(fèi)端的速率天然不匹配。比如訂單瞬間暴漲時(shí)數(shù)據(jù)庫只能每秒寫200條但上游每秒能產(chǎn)生2000個(gè)事件這時(shí)候隊(duì)列就是緩沖水池讓消費(fèi)端穩(wěn)定的200條/秒慢慢消化而不是被瞬時(shí)流量沖垮。實(shí)際業(yè)務(wù)中它帶來的核心價(jià)值是四個(gè)解耦生產(chǎn)者完全不需要知道下游有幾個(gè)消費(fèi)者、每個(gè)消費(fèi)者怎么處理緩沖削峰瞬時(shí)流量壓進(jìn)來時(shí)消息先落在隊(duì)列里消費(fèi)端按自己的節(jié)奏處理可用性提升消費(fèi)端掛掉或者處理失敗時(shí)生產(chǎn)端可以繼續(xù)投遞配合重試機(jī)制能兜底天然支持并行同一個(gè)隊(duì)列可以掛多個(gè)消費(fèi)者這是后續(xù)并行任務(wù)調(diào)度的基礎(chǔ)1.2 三種經(jīng)典實(shí)現(xiàn)方式對(duì)比實(shí)現(xiàn)生產(chǎn)者-消費(fèi)者模式Java 里大體有三條路synchronized wait/notify、Lock Condition、BlockingQueue。初學(xué)者最容易被前兩種繞暈我直接給對(duì)比表。實(shí)現(xiàn)方式同步機(jī)制典型代碼量適用場(chǎng)景必須注意的坑synchronized wait/notify對(duì)象鎖 等待/喚醒多學(xué)習(xí)、極端簡(jiǎn)單場(chǎng)景wait 必須放 while 循環(huán)里必須 notifyAllLock Condition多個(gè)條件隊(duì)列中需要精確喚醒、控制復(fù)雜狀態(tài)容易忘記 unlock多條件容易用錯(cuò)BlockingQueue內(nèi)部同步隊(duì)列最少大部分生產(chǎn)項(xiàng)目選錯(cuò)實(shí)現(xiàn)類容量不設(shè)上限內(nèi)存溢出先看第一版最原始的寫法網(wǎng)上教材里最常見public class BasicProducerConsumer { private final QueueString queue new LinkedList(); private static final int CAPACITY 1000; public synchronized void produce(String item) throws InterruptedException { while (queue.size() CAPACITY) { wait(); // 隊(duì)列滿了暫時(shí)停下等待消費(fèi)端取走 } queue.offer(item); notifyAll(); // 喚醒所有等待的消費(fèi)者 } public synchronized String consume() throws InterruptedException { while (queue.isEmpty()) { wait(); // 隊(duì)列空了等待生產(chǎn)端放入數(shù)據(jù) } String item queue.poll(); notifyAll(); // 喚醒生產(chǎn)者 return item; } }這段代碼看起來沒問題但這里有幾個(gè)隱藏的細(xì)節(jié)。第一個(gè)wait 為什么必須放在 while 而不是 if 里。因?yàn)榇嬖凇疤摷賳拘选本€程可能在沒有被 notify 的情況下醒過來或者被 notifyAll 喚醒后又發(fā)現(xiàn)條件不滿足。用 if 只會(huì)判斷一次醒來直接往下走隊(duì)列還是滿的或空的就會(huì)出問題。while 是二次確認(rèn)條件這是長(zhǎng)期實(shí)踐留下的鐵律。第二個(gè)為什么用 notifyAll 而不是 notify。notify 只隨機(jī)喚醒一個(gè)線程如果喚醒的是一個(gè)同類線程比如生產(chǎn)者喚醒的是另一個(gè)生產(chǎn)者可能永遠(yuǎn)等不到條件滿足。notifyAll 把等待線程全部喚醒讓它們自己再判斷一輪雖然開銷大一點(diǎn)但不會(huì)丟信號(hào)。第三個(gè)鎖對(duì)象必須一致。produce 和 consume 都用 synchronized 鎖了 this如果有一處不小心鎖了另一個(gè)對(duì)象生產(chǎn)者和消費(fèi)者各自的鎖就完全隔離開了整體就是廢的。1.3 為什么很多項(xiàng)目寫著寫著變成了“偽隊(duì)列”這個(gè)標(biāo)題里的“偽隊(duì)列”是我自己起的名字指的是看起來用了 BlockingQueue實(shí)際上既沒解決生產(chǎn)消費(fèi)平衡也沒解決并行調(diào)度只是把一個(gè) List 換了個(gè)線程安全的殼。我自己見過、也親手寫過幾種典型問題。隊(duì)列容量不設(shè)上限。很多人直接new LinkedBlockingQueue()默認(rèn)容量是 Integer.MAX_VALUE等于是無限隊(duì)列。生產(chǎn)端偶爾抽風(fēng)或者流量突增隊(duì)列就能吞掉幾百萬條消息內(nèi)存直接頂著走。這種代碼在測(cè)試環(huán)境永遠(yuǎn)測(cè)不出問題因?yàn)闇y(cè)試流量太小了一上生產(chǎn)就 OOM。消費(fèi)失敗就 continue。取了一條消息process 方法拋了個(gè)異常外層循環(huán)直接 continue消息就永遠(yuǎn)沒了。這在通知推送場(chǎng)景下意味著用戶收不到短信在訂單場(chǎng)景下意味著直接丟單。用線程池但只提交了一個(gè)任務(wù)。啟動(dòng)消費(fèi)者線程池時(shí)寫了個(gè) for 循環(huán)結(jié)果循環(huán)條件寫錯(cuò)了只 submit 了一次后面的人看代碼完全沒察覺只能通過日志里消費(fèi)者的編號(hào)永遠(yuǎn)是 1 來發(fā)現(xiàn)。這節(jié)說這么多不是為了嚇人而是想強(qiáng)調(diào)生產(chǎn)者-消費(fèi)者模式本身不難難的是把它寫成一個(gè)可以長(zhǎng)期迭代、能排查問題、性能不拖后腿的工程組件。后面的并行任務(wù)調(diào)度、注釋優(yōu)化都是圍繞這個(gè)目標(biāo)展開的。2. 并行任務(wù)調(diào)度讓多個(gè)消費(fèi)者真正“并行”起來2.1 單消費(fèi)者瓶頸一個(gè)工人干所有活很多第一次優(yōu)化這套鏈路的人第一反應(yīng)是“把隊(duì)列換得更快一點(diǎn)”“把 LinkedBlockingQueue 換成 ConcurrentLinkedQueue”但實(shí)際上瓶頸往往根本不在隊(duì)列本身而在你只開了一個(gè)消費(fèi)線程。我舉個(gè)例子。假設(shè)一個(gè)批處理任務(wù)每次從隊(duì)列取 1000 條數(shù)據(jù)寫一次數(shù)據(jù)庫耗時(shí)大約 200ms。如果你只啟動(dòng)了一個(gè)消費(fèi)者哪怕隊(duì)列里已經(jīng)堆了 100 萬條任務(wù)每秒鐘能處理的也就是 5000 條吞吐完全被單線程鎖死。而內(nèi)存里明明有 8 核 CPU一個(gè)線程只能占滿一個(gè)核其他核全在空轉(zhuǎn)。單消費(fèi)者模型還有一個(gè)隱性缺點(diǎn)如果消費(fèi)者在處理一條任務(wù)時(shí)偶發(fā)慢調(diào)用比如數(shù)據(jù)庫連接池滿了等待 30 秒那么整個(gè)消費(fèi)鏈路都被卡住隊(duì)列不斷堆積。雖然多消費(fèi)者也會(huì)遇到慢調(diào)用但至少其余消費(fèi)者還能繼續(xù)干活系統(tǒng)不會(huì)“單點(diǎn)停滯”。所以并行任務(wù)調(diào)度要解決的第一件事就是把“一個(gè)消費(fèi)者”變成“一組消費(fèi)者”讓它們并行地消費(fèi)同一個(gè)隊(duì)列。2.2 多消費(fèi)者線程池設(shè)計(jì)與參數(shù)選擇并行消費(fèi)最直接的實(shí)現(xiàn)方式是用一個(gè)線程池來跑多個(gè)消費(fèi)循環(huán)。這里我直接給出一個(gè)在生產(chǎn)環(huán)境穩(wěn)定運(yùn)行過的配置思路而不是甩給你一串神秘參數(shù)。int processors Runtime.getRuntime().availableProcessors(); int consumers Math.max(2, processors); // 保守一點(diǎn)至少 2 個(gè)最多不超過核數(shù)太多 ThreadPoolExecutor consumerPool new ThreadPoolExecutor( consumers, consumers, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue(consumers), // 工作隊(duì)列只存任務(wù)不存業(yè)務(wù)數(shù)據(jù) new NamedThreadFactory(batch-consumer), new ThreadPoolExecutor.CallerRunsPolicy() );這里有幾個(gè)關(guān)鍵點(diǎn)逐一解釋。核心線程數(shù)和最大線程數(shù)設(shè)成一樣都是 consumers。消費(fèi)循環(huán)是常駐任務(wù)不像普通請(qǐng)求那樣有高峰低谷所以不要搞“核心 2 個(gè)、最大 20 個(gè)”這種彈性配置。最大線程數(shù)只在核心線程不夠用時(shí)臨時(shí)擴(kuò)容但消費(fèi)循環(huán)本身是無限阻塞的擴(kuò)容上來的線程出來后立刻又去隊(duì)列里阻塞等待反而可能造成資源浪費(fèi)。固定成一樣行為更可預(yù)測(cè)。工作隊(duì)列的大小設(shè)成 consumers而不是設(shè)成業(yè)務(wù)隊(duì)列的大小。這里很多人的誤區(qū)是把線程池的隊(duì)列當(dāng)成業(yè)務(wù)隊(duì)列用實(shí)際上消費(fèi)線程池的工作隊(duì)列只是存放“消費(fèi)循環(huán)任務(wù)”的任務(wù)數(shù)量很少不需要留大空間。拒絕策略選 CallerRunsPolicy。如果消費(fèi)者線程全部掛掉或者提交任務(wù)失敗CallerRunsPolicy 會(huì)直接在提交線程也就是啟動(dòng)消費(fèi)者的主線程上繼續(xù)執(zhí)行保證任務(wù)不會(huì)無聲無息地丟。相比默認(rèn)的 AbortPolicy 直接拋 RejectedExecutionException這更安全。線程工廠必須自定義命名。用過Executors.defaultThreadFactory()的人都知道線程名一堆 “pool-2-thread-1”出了問題你連是哪個(gè)池的線程都看不出來。用 NamedThreadFactory 可以統(tǒng)一命名成batch-consumer-1、batch-consumer-2日志排查方便非常多。線程數(shù)的估算我一般用這個(gè)公式作為起點(diǎn)N CPU 核數(shù) * (1 等待時(shí)間 / 計(jì)算時(shí)間)。如果任務(wù)是純 IO 型比如批量寫庫、調(diào)第三方接口等待時(shí)間遠(yuǎn)大于計(jì)算時(shí)間線程數(shù)可以放寬到核數(shù)的幾倍甚至幾十倍。如果是 CPU 密集型的解析、加密等任務(wù)線程數(shù)接近核數(shù)就行了開太多反而因?yàn)榫€程切換拖慢速度。實(shí)際場(chǎng)景里我曾經(jīng)把一個(gè)單消費(fèi)者的批量訂單處理改成 8 個(gè)消費(fèi)者并發(fā)處理訂單表從每秒 5000 條提升到 35000 條左右數(shù)據(jù)庫本身成了瓶頸但吞吐的提升是肉眼可見的。這也說明很多并發(fā)優(yōu)化其實(shí)不需要復(fù)雜的算法先把單消費(fèi)者變成多消費(fèi)者效果就立竿見影。2.3 批量合并提交減少鎖競(jìng)爭(zhēng)提升吞吐當(dāng)你已經(jīng)開了多消費(fèi)者下一步值得做的優(yōu)化是批量合并提交。這里的“批量”不是指消費(fèi)者一次只取一條消息而是指攢夠一批再處理尤其適合寫庫、推送等 IO 型操作。我寫一個(gè)實(shí)際可用的消費(fèi)循環(huán)模板注意看 drainTo 的用法public void consumeLoop() { while (!Thread.currentThread().isInterrupted()) { try { ListTask batch new ArrayList(BATCH_SIZE); Task first pendingQueue.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS); if (first null) { continue; // 超時(shí)沒有任務(wù)繼續(xù)下一次循環(huán) } batch.add(first); pendingQueue.drainTo(batch, BATCH_SIZE - 1); // 盡可能多取一些 processBatch(batch); // 批量處理 } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } catch (Exception e) { // 生產(chǎn)環(huán)境這里必須做失敗隔離、重試、死信等兜底 log.error(consume batch failed, e); } } }為什么用 poll 而不是 taketake 會(huì)無限阻塞直到有數(shù)據(jù)如果消費(fèi)端需要優(yōu)雅關(guān)閉線程會(huì)因?yàn)闊o法響應(yīng)中斷而卡住。poll 帶超時(shí)超時(shí)后循環(huán)體有機(jī)會(huì)檢查線程中斷狀態(tài)并退出這是順手解決了一個(gè)很隱晦的停機(jī)問題。為什么用 drainTo 而不是循環(huán)里逐個(gè) take逐個(gè)拿一條處理一條每次都要競(jìng)爭(zhēng)隊(duì)列的鎖。drainTo 一次性把當(dāng)前隊(duì)列里最多 N 條拿出來只需要競(jìng)爭(zhēng)一次鎖批量處理時(shí)鎖開銷大幅降低。之前壓測(cè)過批量 100 條和逐條處理相比整體吞吐能提升 30% 到 50%隊(duì)列越大效果越明顯。還有一個(gè)不能忽略的細(xì)節(jié)batch 的最后一批數(shù)據(jù)要能及時(shí)刷出去不能等積累到 BATCH_SIZE 才處理。上面的代碼里poll 的超時(shí)起到了“兜底”的作用即使一直沒湊夠一批超時(shí)后也會(huì)把已有的幾條拿出去處理避免數(shù)據(jù)無限積壓。2.4 虛擬線程JDK21 下的簡(jiǎn)潔并行方案JDK21 正式發(fā)布虛擬線程之后并行任務(wù)調(diào)度的代碼可以大幅簡(jiǎn)化。虛擬線程最大的特點(diǎn)是“一個(gè)任務(wù)一個(gè)線程”線程本身非常輕量不再需要手動(dòng)設(shè)計(jì)“線程池 消費(fèi)者循環(huán)”這種結(jié)構(gòu)。你只需要這樣寫ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); while (!executor.isShutdown()) { Task task pendingTaskQueue.poll(500, TimeUnit.MILLISECONDS); if (task null) { continue; } executor.submit(() - process(task)); // 每個(gè)任務(wù)一個(gè)虛擬線程 }每個(gè)任務(wù)來了虛擬線程池就創(chuàng)建一個(gè)虛擬線程去執(zhí)行任務(wù)執(zhí)行完虛擬線程自動(dòng)銷毀。因?yàn)闆]有平臺(tái)線程的 1:1 映射百萬級(jí)并發(fā)線程也不會(huì)把內(nèi)存打爆。這個(gè)模型天然支持大批量阻塞 IO 任務(wù)代碼比手動(dòng)維護(hù)線程池簡(jiǎn)單很多。但虛擬線程不是萬能藥。我踩過的幾個(gè)邊界得提前說清楚。CPU 密集型任務(wù)不適合虛擬線程。如果任務(wù)是純計(jì)算比如 JSON 解析、加解密、圖像縮放CPU 核數(shù)有限虛擬線程照樣排隊(duì)等 CPU優(yōu)勢(shì)體現(xiàn)不出來反而因?yàn)檎{(diào)度開銷增加一點(diǎn)性能損失。synchronized 鎖釘住虛擬線程的問題。在虛擬線程在 synchronized 塊內(nèi)阻塞時(shí)會(huì)占住底層的平臺(tái)線程導(dǎo)致平臺(tái)線程被“釘住”大量這種操作會(huì)耗盡載體線程。雖然 JDK24 對(duì) synchronized 做了改進(jìn)但如果你還在跑 JDK21/22遇到阻塞 IO 的場(chǎng)景最好用 ReentrantLock 而不是 synchronized。虛擬線程適合“大量任務(wù) 偶發(fā)阻塞”的場(chǎng)景并不適合“少量常駐任務(wù) 高頻切換”。對(duì)一個(gè)已經(jīng)跑穩(wěn)定的多消費(fèi)者線程池沒有必要為了用虛擬線程而重寫量力而行。3. 代碼注釋的“做得少”與“說清楚”項(xiàng)目的標(biāo)題里特別提到“更簡(jiǎn)潔的注釋和每項(xiàng)改進(jìn)的詳細(xì)解釋”這其實(shí)是我在這個(gè)項(xiàng)目里最有感觸的部分。很多人對(duì)注釋的理解停留在“每行代碼都寫注釋”結(jié)果是代碼上貼滿了廢話真正需要說明的決策理由卻一個(gè)字沒有。3.1 差注釋長(zhǎng)什么樣逐行翻譯式、廢話式我見過太多這種注釋了// 獲取隊(duì)列中的數(shù)據(jù) Task task queue.take(); // 有數(shù)據(jù)就取沒有就一直等 // 處理任務(wù) process(task); // 判斷是否成功 if (task.isSuccess()) { // 記錄日志 log.info(task success); }這段注釋的唯一作用就是占行數(shù)。queue.take()這個(gè)方法名已經(jīng)說明了它在做什么讀者要的不是“它做了什么”而是“為什么在這里用 take 而不是 poll”“為什么沒有設(shè)置超時(shí)”“隊(duì)列空的時(shí)候會(huì)怎樣”。還有一種自我陶醉式的注釋每個(gè)類都要寫作者、創(chuàng)建時(shí)間、修改人/** * author zhangsan * date 2023-05-20 * version 1.0 */這類信息在 Git 提交記錄里本來就有寫在代碼里除了增加維護(hù)負(fù)擔(dān)沒有任何價(jià)值。如果哪天代碼被改了幾十輪作者名字還掛在那里新人以為出錯(cuò)可以找這個(gè)人非常誤導(dǎo)。3.2 用代碼自解釋替代注釋刪注釋容易但刪了之后代碼必須自己說話。我總結(jié)了一套“注釋精簡(jiǎn)三板斧”。第一命名要具體。queue改成pendingTaskQueueprocess改成dispatchTaskdata改成orderEvent。命名精確之后一半注釋都可以刪掉。第二消滅魔法數(shù)。poll(500)里的 500 是什么是超時(shí)毫秒數(shù)。聲明成常量后這個(gè)數(shù)字就有了語義。private static final long POLL_TIMEOUT_MS 500L;第三把復(fù)雜條件提取成方法。代碼里的if (error ! null error.retryCount 3 !isShutdown)很難懂但提取成一個(gè)方法之后if (canRetry(error)) { ... }方法的命名本身就解釋了這段判斷的意圖不需要再寫注釋。這三板斧執(zhí)行完之后你會(huì)驚喜地發(fā)現(xiàn)代碼里還剩的注釋幾乎都是真正值得寫的內(nèi)容。3.3 關(guān)鍵注釋才值得寫原因型注釋、并發(fā)約定、參數(shù)范圍留下來的注釋應(yīng)該長(zhǎng)什么樣核心就一句話注釋只回答“為什么不按常規(guī)來”和“這里有什么約束”。拿前面的消費(fèi)循環(huán)舉例我最終會(huì)在代碼里保留這幾條原因型注釋// 必須用帶超時(shí)的poll否則優(yōu)雅停機(jī)時(shí)線程無法響應(yīng)中斷永遠(yuǎn)卡在take上 Task first pendingQueue.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS); // drainTo一次拿走批量任務(wù)減少鎖競(jìng)爭(zhēng)不能逐個(gè)take否則性能至少下降30% pendingQueue.drainTo(batch, BATCH_SIZE - 1);這兩條注釋都是在讀者可能產(chǎn)生疑問的地方提前解釋而不是解釋代碼本身??吹竭@些注釋的人不用重新踩一遍坑就能知道為什么這么寫。并發(fā)場(chǎng)景下有幾類注釋其實(shí)和代碼邏輯同樣重要不可變約束、可見性約定、鎖順序約定。比如// 該集合只能由consumer線程修改其他線程只讀所以不需要同步 private final MapString, Task cache new ConcurrentHashMap();這種注釋解釋了“為什么不需要同步”“誰來保證線程安全”對(duì)于維護(hù)并發(fā)代碼的人至關(guān)重要。我自己見過太多因?yàn)闆]有寫這類約束后來人看到集合就以為不安全親手加上了一段性能很差的全局鎖。參數(shù)范圍注釋也挺值錢。比如隊(duì)列容量為什么是 2048可以在常量旁邊寫明計(jì)算依據(jù)// 容量 峰值生產(chǎn)速率(500/s) * 單次消費(fèi)延遲(2s) 80%緩沖區(qū)余量 ≈ 1800取整到2048 private static final int QUEUE_CAPACITY 2048;讀者一看就知道改參數(shù)時(shí)從哪里入手而不是拍腦袋調(diào)一個(gè) 99999。3.4 注釋模板的坑現(xiàn)在 IDE 都很智能自動(dòng)生成注釋模板特別方便。我見過有人對(duì)每個(gè)方法都自動(dòng)生成 Javadoc里面全是param xxx 參數(shù)xxxreturn 返回值有的還帶deprecated但沒寫替代方案。這類模板注釋的危險(xiǎn)在于它制造了一種“這個(gè)類很完整很規(guī)范”的假象實(shí)際上讀者得不到任何有效信息。真正要讀代碼的人只能忽略這些模板逐行看代碼。注釋越多噪音越大真話越容易被淹沒。我的習(xí)慣是方法的注釋只在這段邏輯有協(xié)作時(shí)序、并發(fā)約束、事務(wù)邊界時(shí)才寫。字段注釋只在字段含義容易被誤解時(shí)才寫。單行注釋只用于解釋“反直覺”的決策。換句話說注釋的目標(biāo)不是讓代碼顯得很多而是讓代碼顯得很少。如果一段代碼本身已經(jīng)夠簡(jiǎn)潔、命名夠準(zhǔn)確不做注釋是完全正常的。4. 實(shí)例復(fù)盤從單消費(fèi)者基礎(chǔ)版到高性能并行版空講理論沒意思我把一個(gè)簡(jiǎn)化版的實(shí)例完整貼出來。這是從十幾萬行生產(chǎn)代碼里抽出來的模板去掉業(yè)務(wù)細(xì)節(jié)之后大概長(zhǎng)這樣但核心思路和注釋方式都是實(shí)際用過的。4.1 第一版synchronized 基礎(chǔ)版public class BasicPipeline { private final QueueString queue new LinkedList(); private static final int CAPACITY 1000; public synchronized void produce(String item) throws InterruptedException { while (queue.size() CAPACITY) { wait(); } queue.offer(item); notifyAll(); } public synchronized String consume() throws InterruptedException { while (queue.isEmpty()) { wait(); } String item queue.poll(); notifyAll(); return item; } }這一版適合學(xué)習(xí)原理但生產(chǎn)環(huán)境我不會(huì)直接用。原因整個(gè)隊(duì)列只有一個(gè)鎖生產(chǎn)者、消費(fèi)者完全串行互斥隊(duì)列用 LinkedList沒有容量上限保護(hù)這里雖然設(shè)了判斷但如果有多個(gè)生產(chǎn)者并發(fā)判斷size 判斷并不是原子的沒有超時(shí)機(jī)制線程可能無限期阻塞。它最大的問題是沒有發(fā)揮多核并行能力。4.2 第二版BlockingQueue 線程池批量并行版public class BatchParallelPipeline { private final BlockingQueueTask pendingTaskQueue new LinkedBlockingQueue(QUEUE_CAPACITY); private final ExecutorService consumerPool; public BatchParallelPipeline(int consumerCount) { consumerPool new ThreadPoolExecutor( consumerCount, consumerCount, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue(consumerCount), new NamedThreadFactory(batch-consumer), new ThreadPoolExecutor.CallerRunsPolicy() ); } public void start() { for (int i 0; i consumerCount; i) { consumerPool.submit(this::consumeLoop); } } private void consumeLoop() { while (!Thread.currentThread().isInterrupted()) { try { ListTask batch new ArrayList(BATCH_SIZE); Task first pendingTaskQueue.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS); if (first null) { continue; // 空閑時(shí)間不做無意義循環(huán) } batch.add(first); pendingTaskQueue.drainTo(batch, BATCH_SIZE - 1); processBatch(batch); // 這里做業(yè)務(wù)處理 } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } catch (Exception e) { log.error(consume batch failed, batchSize batch.size(), e); // 生產(chǎn)環(huán)境必須有死信或重試策略 } } } public void shutdown() { consumerPool.shutdown(); // 先停止接收新消費(fèi)任務(wù) try { consumerPool.awaitTermination(30, TimeUnit.SECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }這版就是項(xiàng)目里真正能夠上生產(chǎn)的骨架。多消費(fèi)者并行消費(fèi)吞吐量成倍增加。使用 BlockingQueue 之后不再需要手寫 wait/notify線程安全性由隊(duì)列內(nèi)部保證。帶超時(shí)的 poll 加上 drainTo 批量取出既解決了停機(jī)問題又降低了鎖競(jìng)爭(zhēng)。有一點(diǎn)要注意第二版的 processBatch 里如果拋出異常不能直接打死消費(fèi)循環(huán)否則這個(gè)消費(fèi)者線程就退出不干活了。生產(chǎn)環(huán)境我的做法是記錄異常把失敗批次寫進(jìn)一個(gè)死信隊(duì)列給重試線程去處理。雖然會(huì)增加一點(diǎn)復(fù)雜度但數(shù)據(jù)不丟才是最底線的事。4.3 第三版虛擬線程簡(jiǎn)化版JDK21 及以上環(huán)境可以考慮虛擬線程方案。它省去了手動(dòng)管理線程池的環(huán)節(jié)整個(gè)消費(fèi)模型變得更加直接。public void startWithVirtualThreads() { ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); int processors Runtime.getRuntime().availableProcessors(); for (int i 0; i processors; i) { executor.submit(this::consumeLoopV2); } } private void consumeLoopV2() { while (!Thread.currentThread().isInterrupted()) { try { Task task pendingTaskQueue.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS); if (task ! null) { process(task); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }這個(gè)方案有一個(gè)明顯的變化不再需要 drainTo 批量合并了。因?yàn)樘摂M線程足夠輕量為每個(gè)任務(wù)創(chuàng)建一個(gè)虛擬線程的開銷比平臺(tái)線程小得多所以你可以退回到更簡(jiǎn)單的單條處理模型。實(shí)測(cè)下來在 IO 密集型任務(wù)上簡(jiǎn)化版和批量版的吞吐差距可以接受但代碼讀起來通順很多。這個(gè)方案最值得警惕的還是 CPU 密集型任務(wù)。虛擬線程的調(diào)度靠 JVM底層平臺(tái)線程數(shù)量有限一旦所有虛擬線程都在做 CPU 計(jì)算而不是阻塞等待載入過重后整體性能反而下降。用它之前先確認(rèn)你的任務(wù)是不是真的以阻塞 IO 為主。4.4 注釋優(yōu)化前后對(duì)比最后用一個(gè)具體的注釋優(yōu)化對(duì)比收束這節(jié)。這是我從一個(gè)真實(shí)項(xiàng)目里摘出來的重構(gòu)前后對(duì)照。重構(gòu)前// 取出一個(gè)任務(wù) Task task queue.take(); // 執(zhí)行任務(wù) execute(task); // 檢查結(jié)果如果失敗就重試 if (!task.isSuccess()) { retry(task); }重構(gòu)后Task task pendingTaskQueue.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS); // 帶超時(shí)poll避免隊(duì)列空時(shí)線程永遠(yuǎn)阻塞無法優(yōu)雅退出tasknull說明隊(duì)列暫時(shí)為空直接進(jìn)入下一輪 if (task null) { continue; } dispatch(task); // dispatch內(nèi)部已做失敗分類只有可重試異常才走重試隊(duì)列重構(gòu)后的代碼第一眼看上去好像“注釋變少了”實(shí)際上有效信息變多了。讀者一下就知道為什么用 poll 而不用 take、什么情況算暫時(shí)為空、重試邏輯被封裝在哪里。這種注釋風(fēng)格的核心邏輯是把代碼里每一個(gè)“容易讓人困惑的決策”變成明面上的解釋而不是逐行翻譯代碼。注釋的價(jià)值密度高了很多。5. 常見問題與排查技巧實(shí)錄最后這部分全是這些年實(shí)際踩坑后才總結(jié)出來的內(nèi)容。建議收藏出了問題來對(duì)照著看。5.1 死鎖程序卡住日志還打著就是不動(dòng)典型癥狀隊(duì)列消費(fèi)速率降為 0線程池看著還有線程存活但日志不再輸出。排查方式第一時(shí)間打 jstack 抓線程??淳€程處于什么狀態(tài)、卡在哪一行。常見原因用了 synchronized 包住 BlockingQueue 的 take但喚醒條件不對(duì)互相等待多個(gè)任務(wù)持有多個(gè)鎖嵌套獲取兩個(gè)線程互相等對(duì)方釋放消費(fèi)任務(wù)里又調(diào)用了生產(chǎn)者端的方法循環(huán)等待我的經(jīng)驗(yàn)死鎖問題 90% 是鎖的嵌套獲取導(dǎo)致的另外 10% 是 wait/notify 信號(hào)丟失。設(shè)計(jì)上盡量避免在持有鎖的情況下再獲取其他鎖非要嵌套時(shí)必須保證全局鎖順序一致。5.2 隊(duì)列滿了或者空了處理策略是什么隊(duì)列滿的時(shí)候生產(chǎn)者會(huì)阻塞這是 BlockingQueue 的默認(rèn)行為。阻塞的好處是背壓自然傳遞到上游壞處是如果一直滿會(huì)拖慢整個(gè)生產(chǎn)鏈路。我通常給生產(chǎn)者的寫入方法加上超時(shí)boolean offered pendingTaskQueue.offer(task, 1, TimeUnit.SECONDS); if (!offered) { // 隊(duì)列已滿可以選擇走降級(jí)策略丟棄、落本地文件、或者同步處理 }隊(duì)列空的時(shí)候消費(fèi)端循環(huán)需要避免空轉(zhuǎn)。poll 超時(shí)設(shè)為 50ms 到 500ms 比較合適太短會(huì)空轉(zhuǎn)浪費(fèi) CPU太長(zhǎng)會(huì)降低延遲敏感度。如果線程需要優(yōu)雅退出poll 超時(shí)給你一個(gè)定期檢查中斷狀態(tài)的機(jī)會(huì)。5.3 數(shù)據(jù)重復(fù)消費(fèi)或丟失這是并行消費(fèi)最怕的兩件事。數(shù)據(jù)丟失常見原因消費(fèi)線程取出了任務(wù)process 拋異常外層直接 catch 吞掉或者 process 先刪了源數(shù)據(jù)然后自己執(zhí)行失敗數(shù)據(jù)沒了。數(shù)據(jù)重復(fù)常見原因消費(fèi)成功后ack 確認(rèn)失敗消息重新投遞或者 process 執(zhí)行到一半消費(fèi)者宕機(jī)任務(wù)被重新消費(fèi)。應(yīng)對(duì)方案只有兩層第一層消費(fèi)端必須做冪等用唯一業(yè)務(wù)鍵去重第二層處理失敗的消息必須進(jìn)入死信隊(duì)列不能直接丟棄。對(duì)于“先處理還是先確認(rèn)”我的習(xí)慣是先持久化處理結(jié)果再確認(rèn)消費(fèi)這樣即使確認(rèn)失敗也只是重復(fù)處理不會(huì)丟。5.4 線程池線程耗盡任務(wù)全排隊(duì)沒線程干活這個(gè)問題常見于把線程池用于“異步執(zhí)行”業(yè)務(wù)任務(wù)時(shí)。業(yè)務(wù)任務(wù)本身會(huì)去調(diào)慢接口、等待鎖一個(gè)任務(wù)卡住線程池里的線程全被占住新的任務(wù)排在工作隊(duì)列里但永遠(yuǎn)輪不到執(zhí)行。排查方式用 ThreadPoolExecutor 提供的 getActiveCount/getQueue 大小做監(jiān)控出現(xiàn)活躍線程數(shù)長(zhǎng)時(shí)間等于最大線程數(shù)就要警惕了解決方式把任務(wù)按類型拆分成多個(gè)獨(dú)立線程池互相不拖累線程池中的任務(wù)盡量設(shè)置超時(shí)避免永久阻塞情況允許時(shí)考慮用虛擬線程池替代讓阻塞任務(wù)不再占用線程5.5 性能壓測(cè)與調(diào)優(yōu)的一點(diǎn)實(shí)操經(jīng)驗(yàn)最后聊一下怎么驗(yàn)證你的優(yōu)化真的有效。我習(xí)慣的做法是寫一個(gè)小的壓測(cè)入口模擬生產(chǎn)者以不同的速率注入任務(wù)然后統(tǒng)計(jì)消費(fèi)延遲和吞吐long start System.nanoTime(); // 注入100萬條任務(wù) pipeline.produceBatch(1_000_000); long end System.nanoTime(); System.out.println(throughput: 1_000_000L * 1_000_000_000 / (end - start) tasks/s);注意一定要測(cè) P99 延遲而不是只測(cè)平均延遲。并發(fā)場(chǎng)景下平均延遲往往被大多數(shù)快的任務(wù)拉低真正影響用戶體驗(yàn)的是那些卡在尾部的最慢任務(wù)。以前我優(yōu)化完只看平均延遲然后上線后還是被用戶投訴后來加了 P99 監(jiān)控才定位到是某類大任務(wù)偶爾耗時(shí)特別長(zhǎng)把線程池占滿了。調(diào)優(yōu)的順序應(yīng)該是先確認(rèn)瓶頸在 CPU、IO 還是鎖競(jìng)爭(zhēng)上再動(dòng)手改參數(shù)。實(shí)踐里很多問題不是隊(duì)列太小而是消費(fèi)邏輯里有慢查詢也不是線程數(shù)不夠而是線程被無意義的自旋浪費(fèi)了。性能調(diào)優(yōu)最忌諱一上來就盲目改并發(fā)數(shù)連監(jiān)控?cái)?shù)據(jù)都沒看改了一天方向全錯(cuò)了。我個(gè)人的體會(huì)是并發(fā)編程里 70% 的價(jià)值來自于把模型劃分正確剩下的 30% 才來自參數(shù)調(diào)優(yōu)。生產(chǎn)者-消費(fèi)者模式是模型線程池和虛擬線程是模型落地的手段而注釋就是你留給下一個(gè)維護(hù)者的使用說明書。每次重構(gòu)完一套并發(fā)代碼如果能順手記錄這次改動(dòng)的決策理由后面排查問題的成本會(huì)低很多。這個(gè)習(xí)慣堅(jiān)持一年你大概率會(huì)回來感謝自己。