存隊列堵死引發(fā)的性能事故復盤)
事情是從一次大促前壓測開始的。群里突然有人甩了張截圖說訂單支付回調(diào)接口的RT從平時的50毫秒一路飆到了快兩秒后臺還有一堆訂單狀態(tài)卡在“支付中”不更新用戶開始陸續(xù)收到重復的短信提醒。我當時第一反應是數(shù)據(jù)庫又抖了結(jié)果查了一圈發(fā)現(xiàn)問題源頭居然是一個我們自己寫的MessageQueue——更準確點說是我老大當初“順手”寫的一段隊列消費代碼。這個標題里的“Bug”其實不是那種讓你編譯都過不去的錯誤而是一段看著沒什么毛病的代碼在流量稍微上來一點之后把整條異步鏈路活活拖垮了。這個復盤過程很有價值我把它完整記錄下來給所有在做業(yè)務異步化、內(nèi)存隊列、削峰填谷的團隊一個參考。1. 先還原現(xiàn)場老大的“手藝”和半線上事故1.1 業(yè)務背景為什么需要一個隊列先說業(yè)務背景。我們的訂單模塊在用戶支付成功之后要做一連串的“售后動作”更新訂單狀態(tài)、發(fā)短信通知、加積分、推送消息給運營后臺、再通知倉儲系統(tǒng)開始備貨。最早這套邏輯是同步寫在支付回調(diào)接口里的一個請求進來挨個調(diào)用這些服務全部做完之后才返回結(jié)果給客戶端。同步方案在低峰期沒有大問題但有兩個隱患一是接口耗時被下游系統(tǒng)拖累短信服務一抖動支付回調(diào)就跟著超時二是支付成功瞬間的流量通常有尖峰比如整點秒殺、促銷活動開啟的幾分鐘內(nèi)回調(diào)請求會密集進來如果每個請求都要走一遍外部IO數(shù)據(jù)庫和短信接口都可能被打爆。所以就有了這個異步改造支付回調(diào)只負責把“訂單支付成功”這個事件丟進隊列立刻返回后臺再用消費者線程慢慢處理。這正是MessageQueue在這個場景里最核心的價值——解耦和削峰。當時的實現(xiàn)并不復雜也沒有引入RocketMQ或者Kafka這些重量級組件老大直接在應用內(nèi)存里用ArrayBlockingQueue手寫了一個輕量隊列。理論上只要消費者處理速度夠快這套方案完全能支撐現(xiàn)有業(yè)務量。1.2 老大寫的這段代碼到底長啥樣我后來翻到最初的實現(xiàn)大概長這樣public class OrderNotifyQueue { private static final int CAPACITY 1000; private static final BlockingQueueNotifyTask QUEUE new ArrayBlockingQueue(CAPACITY); private static final ExecutorService PRODUCER_POOL Executors.newFixedThreadPool(20); private static final ExecutorService CONSUMER Executors.newSingleThreadExecutor(); static { CONSUMER.execute(new NotifyWorker()); } public static void submit(NotifyTask task) { PRODUCER_POOL.execute(() - { try { QUEUE.put(task); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); } static class NotifyWorker implements Runnable { Override public void run() { while (true) { try { NotifyTask task QUEUE.take(); process(task); } catch (Exception e) { // 失敗就放回隊頭等下輪再處理 try { QUEUE.put(task); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); } } } } private void process(NotifyTask task) { Order order orderMapper.selectById(task.getOrderId()); // 第一次IO查訂單 if (order ! null) { orderMapper.updateStatus(order.getId()); // 第二次IO更新狀態(tài) smsClient.send(order.getMobile(), 您的訂單已支付成功); // 第三次IO外部短信 } } } }說實話第一眼掃過去這段代碼并不算“丑”用了有界隊列防止無腦堆積用put()保證入隊安全消費端單線程也不會出現(xiàn)并發(fā)修改問題失敗重試也“照顧”到了。但它的問題恰恰就藏在這些“看似合理”的細節(jié)里而且這些問題在小流量下根本暴露不出來。1.3 問題為什么沒在代碼評審時被發(fā)現(xiàn)這不是一句“評審不仔細”就能解釋的。我自己后來復盤發(fā)現(xiàn)這類代碼在CR階段經(jīng)常被放過的原因有三個一是評審的重點放在“功能正確性”上。大家會盯著try-catch寫沒寫、空指針會不會出現(xiàn)、事務有沒有加卻很難靜態(tài)地看出一個“吞吐量不達標”的問題——消費者單條處理、逐條打庫這些屬于性能設計范疇不是代碼審查能一眼識破的。二是沒有壓測環(huán)節(jié)。我們的測試環(huán)境數(shù)據(jù)量小隊列永遠吃不飽生產(chǎn)者一扔消費者馬上就消化了CPU、線程、隊列深度這些指標全部正常導致問題被完全掩蓋。等到大促壓測流量上來才徹底現(xiàn)出原形。三是對“隊列模型”缺少統(tǒng)一的規(guī)范。團隊里沒有明確說過“消息必須批量消費”“重試必須退避”“入隊必須具備超時機制”老大按著他以前寫同步接口的思路來寫異步隊列自然就踩進了這幾個典型陷阱。2. 排查過程從CPU飆高到揪出三宗罪2.1 線上現(xiàn)象接口超時和訂單狀態(tài)卡住壓測暴露出來的現(xiàn)象非常直接。支付回調(diào)接口的可接受響應時間是1秒壓測開始后P99直接飆到2秒以上同時后臺監(jiān)控發(fā)現(xiàn)訂單狀態(tài)更新出現(xiàn)大面積積壓積壓數(shù)量一度到了幾萬條短信發(fā)送網(wǎng)關(guān)那邊則報了一堆限流。整個系統(tǒng)并沒有宕機但到處都在“排隊”一副快要被拖垮的樣子。這種“哪里都在排隊”的現(xiàn)象其實是個很好的排查起點。明明只是異步處理慢怎么連支付回調(diào)接口這種“只負責丟消息”的入口都變慢了這一問就能順著生產(chǎn)鏈路摸到隊列本身。2.2 一步步定位top、jstack、Arthas排查第一步先看機器負載。top -Hp找出CPU占用高的線程發(fā)現(xiàn)有一組線程狀態(tài)很奇怪大量線程處于BLOCKED狀態(tài)堆棧全部卡在java.util.concurrent.ArrayBlockingQueue.put。這說明什么說明有一堆生產(chǎn)者在同一個隊列的入隊操作上互相競爭、排隊等待。第二步用jstack抓線程棧重點看消息隊列相關(guān)的線程和業(yè)務回調(diào)線程。很快就能看到幾十個http-nio-exec-*線程都在put方法上等待鎖而真正干活的消費者線程只有一個它正深度卡在orderMapper.selectById之后的等待里——準確說是在等短信服務響應。單消費者線程本身處理一條消息就要經(jīng)歷三次串行IO平均耗時80到100毫秒算下來吞吐量大概也就每秒12條。第三步我用Arthas做細化定位。thread -n 3再看線程CPU排行。接著用trace命令跟蹤NotifyWorker.process方法發(fā)現(xiàn)耗時大部分落在兩個地方一次訂單查詢SQL平均15毫秒、一次短信發(fā)送外部HTTP調(diào)用平均50毫秒。單條消息就要消耗掉70到80毫秒的處理時間而生產(chǎn)速率在峰值時一秒有三四百條這個隊列不堵才奇怪。2.3 壓測復現(xiàn)與根因確認光看線程棧還不夠必須壓測復現(xiàn)。我們在壓測環(huán)境重新搭了一套一模一樣的代碼用腳本模擬支付回調(diào)按峰值速率每秒400個請求往隊列里灌。10分鐘之后隊列長度直接頂?shù)?000的容量上限隨后所有put請求開始阻塞回調(diào)接口RT應聲上漲。這一步把因果鏈坐實了生產(chǎn)速率大于消費速率隊列開始堆積堆積到容量上限后put無限期等待于是生產(chǎn)者線程——也就是支付回調(diào)的線程——全部被堵在隊列上回調(diào)線程被堵住接口RT自然飆升。整個過程層層傳導跟多米諾骨牌一樣。2.4 三宗罪隊列堵死的核心原因根因確認之后我把問題歸結(jié)為三宗罪第一宗罪入隊用put()無限阻塞。put()的語義是“隊列滿了一直等”這在低流量下沒什么但在峰值時段隊列一旦滿了所有回調(diào)線程都會無限期掛在入隊操作上。最要命的是這種等待沒有超時、沒有降級、沒有熔斷外部流量還在不斷增加線程池所有線程全部被占滿于是接口徹底失去響應能力。第二宗罪消費者單條處理全鏈路串行IO。每消費一條消息就要查一次庫、更新一次狀態(tài)、發(fā)一次短信。這三個操作全是IO型操作單條耗時就接近百毫秒。一個單線程消費者一天最多也就處理一百萬條消息聽著不少但頂不住秒殺瞬間的尖峰流量。隊列消費端的設計完全不符合“批量”這個最基本的優(yōu)化思路。第三宗罪失敗重試直接放回隊頭。代碼里一旦process拋出異常就把任務重新放回隊列頭部。這個操作有兩個問題第一如果某條消息持續(xù)失敗它會堵在隊首一直占著生產(chǎn)者的名額第二put操作本身也可能阻塞失敗的線程會把自己卡在重試入隊上。而且恢復正常順序會被打亂用戶收到的短信順序可能錯亂業(yè)務上非常尷尬。2.5 一個隱藏的“背鍋位”單線程消費者除了上面三條還有一個容易被忽略的設計問題消費線程只有一個。單線程消費天然無法利用多核CPU遇到磁盤IO或者外部接口等待時CPU只能空轉(zhuǎn)。很多團隊看到“單線程處理不存在并發(fā)問題”就覺得安全卻忘了消費能力一樣是硬指標。增加合理的消費線程數(shù)本來就是隊列設計的一部分。3. 優(yōu)化方案每一處改動背后的理由3.1 把“按條消費”改成“攢批消費”第一個改動也是最核心的一個消費者不再一條一條處理而是“攢一批、處理一批”。這里用到了BlockingQueue的drainTo方法配合超時輪詢static class BatchNotifyWorker implements Runnable { private static final int MAX_BATCH_SIZE 500; private static final long POLL_TIMEOUT_MS 200; Override public void run() { ListNotifyTask batch new ArrayList(MAX_BATCH_SIZE); while (!Thread.currentThread().isInterrupted()) { batch.clear(); // 先嘗試拿到第一條最多等200ms NotifyTask first QUEUE.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS); if (first null) { continue; // 隊列空閑避免空轉(zhuǎn) } batch.add(first); // 把當前隊列里能拿的都拿出來 QUEUE.drainTo(batch, MAX_BATCH_SIZE - 1); processBatch(batch); } } private void processBatch(ListNotifyTask tasks) { ListInteger orderIds tasks.stream().map(NotifyTask::getOrderId).collect(toList()); ListOrder orders orderMapper.selectByIds(orderIds); // 一次批量查詢 orderMapper.batchUpdateStatus(orders); // 一次批量更新 smsClient.batchSend(orders.stream() .map(o - new SmsRequest(o.getMobile(), 您的訂單已支付成功)) .collect(toList())); // 一次合并發(fā)送 } }核心思路是把“三次串行IO”壓縮成“三次批量IO”。原來100條消息要查100次庫、更新100次庫、發(fā)100次短信現(xiàn)在變成1次批量查詢、1次批量更新、1次合并短信發(fā)送。數(shù)據(jù)庫批量操作的耗時并不是按條數(shù)線性增長的100條的批量更新和1條的單條更新相比耗時可能只多了一倍不到但單位吞吐量翻了近百倍。這里有兩點要注意drainTo最大取499條加上前面那一條正好湊滿500而第一次poll設置200毫秒超時是為了在隊列空閑時不頻繁空轉(zhuǎn)同時保證只要隊列里有消息最多等200毫秒就能積攢出一批。3.2 入隊策略put換offer把無限阻塞改成有界等待生產(chǎn)端改動也很關(guān)鍵put()換成offer()加超時超時之后走降級邏輯public static boolean submit(NotifyTask task) { // 隊列滿時最多等200ms再不行就降級 return QUEUE.offer(task, 200, TimeUnit.MILLISECONDS); } // 使用處 boolean accepted OrderNotifyQueue.submit(task); if (!accepted) { // 降級策略寫入本地文件或DB待重發(fā)也可以直接走同步處理 fallbackSave(task); }為什么不是把隊列直接改成無界無界隊列看起來“永遠不會拒絕消息”但代價是內(nèi)存被無限占用最終觸發(fā)FGC甚至OOM。線上的資源是有限的削峰的本質(zhì)是“暫時存不下了就先擋住”而不是“有多少都硬接”。有界隊列加超時本質(zhì)上是一種背壓機制——當下游處理不過來了上游也要感覺到壓力然后想辦法降級或者限流而不是把整個系統(tǒng)拖死。超時時間選200毫秒是和支付回調(diào)的可用性指標對齊的回調(diào)本身還有一系列其他邏輯入隊最多占200毫秒再往上去就觸發(fā)降級保證接口不會被隊列拖到超時。3.3 重試隊列獨立出來用延遲隊列做退避重試邏輯的優(yōu)化核心原則是“失敗的消息不能回主隊列頭必須退避而且不能擠壓新消息”。我改成了獨立的延遲隊列private static final DelayQueueDelayItemNotifyTask RETRY_QUEUE new DelayQueue(); private static final long[] RETRY_DELAYS_MS {5_000, 30_000, 60_000, 300_000}; public static void retryLater(NotifyTask task, int retryCount) { long delay RETRY_DELAYS_MS[Math.min(retryCount, RETRY_DELAYS_MS.length - 1)]; RETRY_QUEUE.put(new DelayItem(task, System.currentTimeMillis() delay)); }消費線程在處理完主隊列任務后會檢查延遲隊列是否有到期任務有就取出來重新放回主隊列。這樣既保證了失敗任務會重試又不會讓它們卡在主隊列里餓死后來的正常消息。分級退避從5秒到300秒逐級拉長避免某條消息持續(xù)失敗的時候瘋狂重試打爆下游。為什么不用ScheduledExecutorService因為延遲隊列天然支持“按到期時間排序取出”而定時線程池更偏向“周期性任務”處理那種“到點觸發(fā)一次、不再管了”的邏輯還行要做“延遲之后重新入隊”這種動態(tài)流轉(zhuǎn)DelayQueue更順手。3.4 消費者線程數(shù)怎么算消費線程從1個改成多個但線程數(shù)不能拍腦袋。我用的估算公式是這個對于IO密集型任務線程數(shù) ≈ CPU核數(shù) × (1 平均等待時間 / 平均計算時間)我們壓測環(huán)境是8核批量處理中數(shù)據(jù)庫批量IO加外部短信調(diào)用的等待時間大約占75%本地CPU計算和JSON序列化大約占25%。算下來等待時間和計算時間的比值大約是3。于是線程數(shù) ≈ 8 × (1 3) 32但這是理論上限實際并沒有直接配32——線程太多會導致數(shù)據(jù)庫連接池和短信網(wǎng)關(guān)連接數(shù)不夠用反而加劇競爭。折中一下我配了4個消費線程每個線程一次處理500條4個線程疊加起來峰值吞吐量已經(jīng)能達到每秒幾百條完全覆蓋當前生產(chǎn)速率。多線程時代的“安全”不是單線程而是“可控數(shù)量的多線程加上批量消費”。3.5 更進一步要不要直接上Disruptor或者專業(yè)MQ優(yōu)化完自家內(nèi)存隊列之后團隊里也有人問都這么費勁了為什么不直接換Kafka或者RocketMQ這個問題我認真想了一下。結(jié)論是業(yè)務場景和成本決定架構(gòu)選型。我們這里只是一個訂單模塊的內(nèi)部異步通知數(shù)據(jù)量級遠沒有到需要分布式消息隊列的水平引入專業(yè)MQ意味著要部署B(yǎng)roker集群、維護Topic、處理分區(qū)和消費者組關(guān)系運維成本和復雜度一下就上來了。Disruptor的核心優(yōu)勢是無鎖環(huán)形隊列適合超高吞吐的場景但對批量處理、重試退避這類業(yè)務邏輯沒有直接幫助反而因為API抽象更底層團隊學習成本更高。所以最終保留了自己優(yōu)化的內(nèi)存隊列但加了一條規(guī)則如果未來積壓量持續(xù)超過內(nèi)存隊列上限的50%或者出現(xiàn)多機房部署需求就切換專業(yè)MQ。4. 壓測效果與上線驗證4.1 壓測方案與測試腳本優(yōu)化完成后不能直接上生產(chǎn)必須先壓測。我們復用了之前那套腳本模擬支付回調(diào)接口每秒產(chǎn)生400個訂單通知任務持續(xù)壓測30分鐘。同時把生產(chǎn)者和消費者的關(guān)鍵指標隊列深度、消費耗時、重試次數(shù)都打印到日志再配合監(jiān)控系統(tǒng)看趨勢。壓測腳本我寫得很簡單核心就是通過一個循環(huán)接口往隊列里塞數(shù)據(jù)再統(tǒng)計每秒成功入隊的數(shù)量和回調(diào)響應時間。這個腳本的價值不在于代碼多精妙而在于它能夠真實逼近線上峰值速率讓問題在“發(fā)布之前”暴露出來。4.2 優(yōu)化前后的核心指標對比壓測結(jié)束后我把數(shù)據(jù)整理成了表格結(jié)論非常直觀指標優(yōu)化前優(yōu)化后說明消費者吞吐量約12條/秒約600條/秒批量消費帶來的數(shù)量級提升支付回調(diào)接口P99耗時2秒180毫秒入隊不再阻塞回調(diào)線程隊列積壓數(shù)打滿1000持續(xù)觸頂峰值不超過200消費能力覆蓋生產(chǎn)速率系統(tǒng)CPU占用90%以上約35%線程等待減少、批量IO減少空轉(zhuǎn)短信發(fā)送調(diào)用次數(shù)每單1次HTTP合并批量發(fā)送下游限流問題基本消失失敗消息重試回隊頭干擾正常消息延遲隊列分級退避不再影響正常消費順序最顯著的變化是吞吐量從12條/秒提升到600條/秒整整50倍。這個結(jié)果并不夸張因為優(yōu)化的核心不是“讓單條消息處理得更快”而是“一次處理一批消息”——單條處理耗時80毫秒批量500條處理耗時也就200毫秒相當于單條均攤時間從80毫秒降到了0.4毫秒。4.3 灰度發(fā)布與監(jiān)控落地指標好看也不能一股腦全量上。我們分了三個步驟灰度先切10%流量觀察半天重點看監(jiān)控面板上的隊列深度、消費延遲、重試次數(shù)三個指標有沒有異常。半天沒問題再放到50%再觀察半天最后才全量。全量之后持續(xù)盯了72小時確認短信發(fā)送量正常、訂單狀態(tài)更新無積壓、回調(diào)接口RT穩(wěn)定這次優(yōu)化才算真正閉環(huán)。監(jiān)控是比壓測更重要的東西。壓測只能驗證“在那個時間點沒問題”線上流量像潮水一樣隨時變化沒有監(jiān)控就只能靠用戶投訴來發(fā)現(xiàn)問題。我后來專門給隊列加了三塊看板隊列實時深度、單條消息從入隊到消費完成的端到端延遲、失敗重試的次數(shù)和級別。有了這三個指標以后隊列再出問題打開監(jiān)控就能定位到是生產(chǎn)太快還是消費太慢或者是重試風暴。4.4 上線后的一點觀察上線后正好趕上一次小規(guī)模促銷業(yè)務流量比平時翻了3倍多系統(tǒng)穩(wěn)如老狗。對比之前壓測就崩的狀態(tài)差異非常明顯。銷售那邊還跑來問“最近短信怎么發(fā)得這么快”我們只能笑笑說“改了個Bug”。5. 常見問題與排坑經(jīng)驗速查5.1 內(nèi)存隊列參數(shù)速查表這次踩坑之后我把內(nèi)存隊列的參數(shù)整理成一張表發(fā)給團隊所有人參考。如果你們也在自己寫內(nèi)存隊列可以直接抄參數(shù)推薦值理由隊列容量按峰值積壓量的2到3倍設置太小容易觸發(fā)背壓太大可能內(nèi)存浪費單批次大小200到500太大時單批處理時間過長太小體現(xiàn)不出批量優(yōu)勢入隊超時100到300毫秒與接口可容忍延遲對齊超時即降級消費線程數(shù)CPU核數(shù)×(1等待/計算時間)再折半防止把DB連接池和下游連接數(shù)打滿失敗重試延遲5秒/30秒/60秒/300秒分級逐級退避避免重試風暴批量poll超時100到200毫秒平衡積攢時間和空轉(zhuǎn)消耗5.2 消息丟失、重復、亂序怎么取舍消息隊列的本質(zhì)是異步和削峰它不可能同時保證“不丟失、不重復、不亂序”這是分布式系統(tǒng)的基本約束。你必須根據(jù)業(yè)務的容忍度做取舍。比如我們訂單狀態(tài)更新這個場景重復通知用戶是難以接受的所以消費者端必須做冪等更新訂單狀態(tài)前先判斷當前狀態(tài)已經(jīng)是“已支付”的就不再重復發(fā)短信短信發(fā)送接口也做了業(yè)務冪等鍵同一個訂單號在短時間內(nèi)不會重復發(fā)送。消息丟失的兜底則是靠降級落庫如果隊列滿了實在入不了隊就把任務寫進本地待辦表由一個定時任務每5分鐘掃一次重新提交進隊列。這樣一來極端情況下最多延遲幾分鐘但不會丟。亂序問題在我們場景里影響不大因為訂單支付成功通知本質(zhì)上不依賴嚴格順序。但如果你的業(yè)務是“先改狀態(tài)再發(fā)短信”這種強順序鏈路那就需要在消息體里帶一個業(yè)務序列號消費端按序列號做排序或者丟棄過期消息而不是單純依賴隊列的有序性。5.3 這次踩坑后總結(jié)的幾條團隊規(guī)約這次事故給團隊帶來的最大收益不是代碼優(yōu)化本身而是幾個流程性的改變第一所有涉及隊列、線程池、異步處理的代碼CR時必須有壓測記錄或者明確的性能指標預估不能只聊邏輯正確性。第二使用內(nèi)存隊列必須遵守“有界超時降級”三件套不允許在生產(chǎn)代碼里出現(xiàn)裸的put()無限阻塞。第三重試邏輯必須和正常消息分離要么用延遲隊列要么用專門的待處理表禁止把失敗任務放回隊列頭部。第四條是我個人加上的——任何異步鏈路上線前必須有隊列深度的監(jiān)控告警和消費延遲的監(jiān)控告警。沒有監(jiān)控就是睜著眼睛把系統(tǒng)交給運氣。5.4 關(guān)于“Bug”這件事的一點感想說到“Bug”網(wǎng)上總能看到類似“codex磁盤bug”或者“winsxs bug”這類聽著就讓人頭大的問題但說實話這些離我們?nèi)粘i_發(fā)太遠了。真正讓我們寢食難安的往往是老大寫的這段“當時看起來沒問題”的隊列代碼它不報錯、不崩潰只是在你最需要它扛住流量的時候安靜地堵在那里把所有入口都堵死。排這種Bug最難的從來不是修復而是找到那個讓所有線索串起來的核心因果鏈——隊列滿了生產(chǎn)者的線程全被卡在入隊上回調(diào)接口才變慢背后是消費端單條處理吞吐不夠。我個人在實際排查和優(yōu)化過程中最大的體會是遇到性能相關(guān)的線上問題永遠不要急著改代碼。先把線程棧抓下來把壓測復現(xiàn)跑出來把數(shù)據(jù)擺在桌面上再動手。因為只有數(shù)據(jù)能告訴你真正該優(yōu)化的是什么——而不是你“感覺”該優(yōu)化的是什么。這個習慣在代碼評審里同樣適用看到一段“能用”的代碼多問一句“流量翻十倍它還扛得住嗎”很多線上事故就根本不會發(fā)生。最后再分享一個實用的小技巧排查完類似問題之后把當時的線程dump、壓測腳本和優(yōu)化對比數(shù)據(jù)留在一個專門的問題追蹤文檔里。下次再遇到隊列或者線程池相關(guān)的性能問題先翻這個文檔大概率能省下半天排查時間。