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