从50ms到2秒:一次内存队列堵死引发的性能事故复盘 事情是从一次大促前压测开始的。群里突然有人甩了张截图说订单支付回调接口的RT从平时的50毫秒一路飙到了快两秒后台还有一堆订单状态卡在“支付中”不更新用户开始陆续收到重复的短信提醒。我当时第一反应是数据库又抖了结果查了一圈发现问题源头居然是一个我们自己写的MessageQueue——更准确点说是我老大当初“顺手”写的一段队列消费代码。这个标题里的“Bug”其实不是那种让你编译都过不去的错误而是一段看着没什么毛病的代码在流量稍微上来一点之后把整条异步链路活活拖垮了。这个复盘过程很有价值我把它完整记录下来给所有在做业务异步化、内存队列、削峰填谷的团队一个参考。1. 先还原现场老大的“手艺”和半线上事故1.1 业务背景为什么需要一个队列先说业务背景。我们的订单模块在用户支付成功之后要做一连串的“售后动作”更新订单状态、发短信通知、加积分、推送消息给运营后台、再通知仓储系统开始备货。最早这套逻辑是同步写在支付回调接口里的一个请求进来挨个调用这些服务全部做完之后才返回结果给客户端。同步方案在低峰期没有大问题但有两个隐患一是接口耗时被下游系统拖累短信服务一抖动支付回调就跟着超时二是支付成功瞬间的流量通常有尖峰比如整点秒杀、促销活动开启的几分钟内回调请求会密集进来如果每个请求都要走一遍外部IO数据库和短信接口都可能被打爆。所以就有了这个异步改造支付回调只负责把“订单支付成功”这个事件丢进队列立刻返回后台再用消费者线程慢慢处理。这正是MessageQueue在这个场景里最核心的价值——解耦和削峰。当时的实现并不复杂也没有引入RocketMQ或者Kafka这些重量级组件老大直接在应用内存里用ArrayBlockingQueue手写了一个轻量队列。理论上只要消费者处理速度够快这套方案完全能支撑现有业务量。1.2 老大写的这段代码到底长啥样我后来翻到最初的实现大概长这样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更新状态 smsClient.send(order.getMobile(), 您的订单已支付成功); // 第三次IO外部短信 } } } }说实话第一眼扫过去这段代码并不算“丑”用了有界队列防止无脑堆积用put()保证入队安全消费端单线程也不会出现并发修改问题失败重试也“照顾”到了。但它的问题恰恰就藏在这些“看似合理”的细节里而且这些问题在小流量下根本暴露不出来。1.3 问题为什么没在代码评审时被发现这不是一句“评审不仔细”就能解释的。我自己后来复盘发现这类代码在CR阶段经常被放过的原因有三个一是评审的重点放在“功能正确性”上。大家会盯着try-catch写没写、空指针会不会出现、事务有没有加却很难静态地看出一个“吞吐量不达标”的问题——消费者单条处理、逐条打库这些属于性能设计范畴不是代码审查能一眼识破的。二是没有压测环节。我们的测试环境数据量小队列永远吃不饱生产者一扔消费者马上就消化了CPU、线程、队列深度这些指标全部正常导致问题被完全掩盖。等到大促压测流量上来才彻底现出原形。三是对“队列模型”缺少统一的规范。团队里没有明确说过“消息必须批量消费”“重试必须退避”“入队必须具备超时机制”老大按着他以前写同步接口的思路来写异步队列自然就踩进了这几个典型陷阱。2. 排查过程从CPU飙高到揪出三宗罪2.1 线上现象接口超时和订单状态卡住压测暴露出来的现象非常直接。支付回调接口的可接受响应时间是1秒压测开始后P99直接飙到2秒以上同时后台监控发现订单状态更新出现大面积积压积压数量一度到了几万条短信发送网关那边则报了一堆限流。整个系统并没有宕机但到处都在“排队”一副快要被拖垮的样子。这种“哪里都在排队”的现象其实是个很好的排查起点。明明只是异步处理慢怎么连支付回调接口这种“只负责丢消息”的入口都变慢了这一问就能顺着生产链路摸到队列本身。2.2 一步步定位top、jstack、Arthas排查第一步先看机器负载。top -Hp找出CPU占用高的线程发现有一组线程状态很奇怪大量线程处于BLOCKED状态堆栈全部卡在java.util.concurrent.ArrayBlockingQueue.put。这说明什么说明有一堆生产者在同一个队列的入队操作上互相竞争、排队等待。第二步用jstack抓线程栈重点看消息队列相关的线程和业务回调线程。很快就能看到几十个http-nio-exec-*线程都在put方法上等待锁而真正干活的消费者线程只有一个它正深度卡在orderMapper.selectById之后的等待里——准确说是在等短信服务响应。单消费者线程本身处理一条消息就要经历三次串行IO平均耗时80到100毫秒算下来吞吐量大概也就每秒12条。第三步我用Arthas做细化定位。thread -n 3再看线程CPU排行。接着用trace命令跟踪NotifyWorker.process方法发现耗时大部分落在两个地方一次订单查询SQL平均15毫秒、一次短信发送外部HTTP调用平均50毫秒。单条消息就要消耗掉70到80毫秒的处理时间而生产速率在峰值时一秒有三四百条这个队列不堵才奇怪。2.3 压测复现与根因确认光看线程栈还不够必须压测复现。我们在压测环境重新搭了一套一模一样的代码用脚本模拟支付回调按峰值速率每秒400个请求往队列里灌。10分钟之后队列长度直接顶到1000的容量上限随后所有put请求开始阻塞回调接口RT应声上涨。这一步把因果链坐实了生产速率大于消费速率队列开始堆积堆积到容量上限后put无限期等待于是生产者线程——也就是支付回调的线程——全部被堵在队列上回调线程被堵住接口RT自然飙升。整个过程层层传导跟多米诺骨牌一样。2.4 三宗罪队列堵死的核心原因根因确认之后我把问题归结为三宗罪第一宗罪入队用put()无限阻塞。put()的语义是“队列满了一直等”这在低流量下没什么但在峰值时段队列一旦满了所有回调线程都会无限期挂在入队操作上。最要命的是这种等待没有超时、没有降级、没有熔断外部流量还在不断增加线程池所有线程全部被占满于是接口彻底失去响应能力。第二宗罪消费者单条处理全链路串行IO。每消费一条消息就要查一次库、更新一次状态、发一次短信。这三个操作全是IO型操作单条耗时就接近百毫秒。一个单线程消费者一天最多也就处理一百万条消息听着不少但顶不住秒杀瞬间的尖峰流量。队列消费端的设计完全不符合“批量”这个最基本的优化思路。第三宗罪失败重试直接放回队头。代码里一旦process抛出异常就把任务重新放回队列头部。这个操作有两个问题第一如果某条消息持续失败它会堵在队首一直占着生产者的名额第二put操作本身也可能阻塞失败的线程会把自己卡在重试入队上。而且恢复正常顺序会被打乱用户收到的短信顺序可能错乱业务上非常尴尬。2.5 一个隐藏的“背锅位”单线程消费者除了上面三条还有一个容易被忽略的设计问题消费线程只有一个。单线程消费天然无法利用多核CPU遇到磁盘IO或者外部接口等待时CPU只能空转。很多团队看到“单线程处理不存在并发问题”就觉得安全却忘了消费能力一样是硬指标。增加合理的消费线程数本来就是队列设计的一部分。3. 优化方案每一处改动背后的理由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; // 队列空闲避免空转 } 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())); // 一次合并发送 } }核心思路是把“三次串行IO”压缩成“三次批量IO”。原来100条消息要查100次库、更新100次库、发100次短信现在变成1次批量查询、1次批量更新、1次合并短信发送。数据库批量操作的耗时并不是按条数线性增长的100条的批量更新和1条的单条更新相比耗时可能只多了一倍不到但单位吞吐量翻了近百倍。这里有两点要注意drainTo最大取499条加上前面那一条正好凑满500而第一次poll设置200毫秒超时是为了在队列空闲时不频繁空转同时保证只要队列里有消息最多等200毫秒就能积攒出一批。3.2 入队策略put换offer把无限阻塞改成有界等待生产端改动也很关键put()换成offer()加超时超时之后走降级逻辑public static boolean submit(NotifyTask task) { // 队列满时最多等200ms再不行就降级 return QUEUE.offer(task, 200, TimeUnit.MILLISECONDS); } // 使用处 boolean accepted OrderNotifyQueue.submit(task); if (!accepted) { // 降级策略写入本地文件或DB待重发也可以直接走同步处理 fallbackSave(task); }为什么不是把队列直接改成无界无界队列看起来“永远不会拒绝消息”但代价是内存被无限占用最终触发FGC甚至OOM。线上的资源是有限的削峰的本质是“暂时存不下了就先挡住”而不是“有多少都硬接”。有界队列加超时本质上是一种背压机制——当下游处理不过来了上游也要感觉到压力然后想办法降级或者限流而不是把整个系统拖死。超时时间选200毫秒是和支付回调的可用性指标对齐的回调本身还有一系列其他逻辑入队最多占200毫秒再往上去就触发降级保证接口不会被队列拖到超时。3.3 重试队列独立出来用延迟队列做退避重试逻辑的优化核心原则是“失败的消息不能回主队列头必须退避而且不能挤压新消息”。我改成了独立的延迟队列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秒逐级拉长避免某条消息持续失败的时候疯狂重试打爆下游。为什么不用ScheduledExecutorService因为延迟队列天然支持“按到期时间排序取出”而定时线程池更偏向“周期性任务”处理那种“到点触发一次、不再管了”的逻辑还行要做“延迟之后重新入队”这种动态流转DelayQueue更顺手。3.4 消费者线程数怎么算消费线程从1个改成多个但线程数不能拍脑袋。我用的估算公式是这个对于IO密集型任务线程数 ≈ CPU核数 × (1 平均等待时间 / 平均计算时间)我们压测环境是8核批量处理中数据库批量IO加外部短信调用的等待时间大约占75%本地CPU计算和JSON序列化大约占25%。算下来等待时间和计算时间的比值大约是3。于是线程数 ≈ 8 × (1 3) 32但这是理论上限实际并没有直接配32——线程太多会导致数据库连接池和短信网关连接数不够用反而加剧竞争。折中一下我配了4个消费线程每个线程一次处理500条4个线程叠加起来峰值吞吐量已经能达到每秒几百条完全覆盖当前生产速率。多线程时代的“安全”不是单线程而是“可控数量的多线程加上批量消费”。3.5 更进一步要不要直接上Disruptor或者专业MQ优化完自家内存队列之后团队里也有人问都这么费劲了为什么不直接换Kafka或者RocketMQ这个问题我认真想了一下。结论是业务场景和成本决定架构选型。我们这里只是一个订单模块的内部异步通知数据量级远没有到需要分布式消息队列的水平引入专业MQ意味着要部署Broker集群、维护Topic、处理分区和消费者组关系运维成本和复杂度一下就上来了。Disruptor的核心优势是无锁环形队列适合超高吞吐的场景但对批量处理、重试退避这类业务逻辑没有直接帮助反而因为API抽象更底层团队学习成本更高。所以最终保留了自己优化的内存队列但加了一条规则如果未来积压量持续超过内存队列上限的50%或者出现多机房部署需求就切换专业MQ。4. 压测效果与上线验证4.1 压测方案与测试脚本优化完成后不能直接上生产必须先压测。我们复用了之前那套脚本模拟支付回调接口每秒产生400个订单通知任务持续压测30分钟。同时把生产者和消费者的关键指标队列深度、消费耗时、重试次数都打印到日志再配合监控系统看趋势。压测脚本我写得很简单核心就是通过一个循环接口往队列里塞数据再统计每秒成功入队的数量和回调响应时间。这个脚本的价值不在于代码多精妙而在于它能够真实逼近线上峰值速率让问题在“发布之前”暴露出来。4.2 优化前后的核心指标对比压测结束后我把数据整理成了表格结论非常直观指标优化前优化后说明消费者吞吐量约12条/秒约600条/秒批量消费带来的数量级提升支付回调接口P99耗时2秒180毫秒入队不再阻塞回调线程队列积压数打满1000持续触顶峰值不超过200消费能力覆盖生产速率系统CPU占用90%以上约35%线程等待减少、批量IO减少空转短信发送调用次数每单1次HTTP合并批量发送下游限流问题基本消失失败消息重试回队头干扰正常消息延迟队列分级退避不再影响正常消费顺序最显著的变化是吞吐量从12条/秒提升到600条/秒整整50倍。这个结果并不夸张因为优化的核心不是“让单条消息处理得更快”而是“一次处理一批消息”——单条处理耗时80毫秒批量500条处理耗时也就200毫秒相当于单条均摊时间从80毫秒降到了0.4毫秒。4.3 灰度发布与监控落地指标好看也不能一股脑全量上。我们分了三个步骤灰度先切10%流量观察半天重点看监控面板上的队列深度、消费延迟、重试次数三个指标有没有异常。半天没问题再放到50%再观察半天最后才全量。全量之后持续盯了72小时确认短信发送量正常、订单状态更新无积压、回调接口RT稳定这次优化才算真正闭环。监控是比压测更重要的东西。压测只能验证“在那个时间点没问题”线上流量像潮水一样随时变化没有监控就只能靠用户投诉来发现问题。我后来专门给队列加了三块看板队列实时深度、单条消息从入队到消费完成的端到端延迟、失败重试的次数和级别。有了这三个指标以后队列再出问题打开监控就能定位到是生产太快还是消费太慢或者是重试风暴。4.4 上线后的一点观察上线后正好赶上一次小规模促销业务流量比平时翻了3倍多系统稳如老狗。对比之前压测就崩的状态差异非常明显。销售那边还跑来问“最近短信怎么发得这么快”我们只能笑笑说“改了个Bug”。5. 常见问题与排坑经验速查5.1 内存队列参数速查表这次踩坑之后我把内存队列的参数整理成一张表发给团队所有人参考。如果你们也在自己写内存队列可以直接抄参数推荐值理由队列容量按峰值积压量的2到3倍设置太小容易触发背压太大可能内存浪费单批次大小200到500太大时单批处理时间过长太小体现不出批量优势入队超时100到300毫秒与接口可容忍延迟对齐超时即降级消费线程数CPU核数×(1等待/计算时间)再折半防止把DB连接池和下游连接数打满失败重试延迟5秒/30秒/60秒/300秒分级逐级退避避免重试风暴批量poll超时100到200毫秒平衡积攒时间和空转消耗5.2 消息丢失、重复、乱序怎么取舍消息队列的本质是异步和削峰它不可能同时保证“不丢失、不重复、不乱序”这是分布式系统的基本约束。你必须根据业务的容忍度做取舍。比如我们订单状态更新这个场景重复通知用户是难以接受的所以消费者端必须做幂等更新订单状态前先判断当前状态已经是“已支付”的就不再重复发短信短信发送接口也做了业务幂等键同一个订单号在短时间内不会重复发送。消息丢失的兜底则是靠降级落库如果队列满了实在入不了队就把任务写进本地待办表由一个定时任务每5分钟扫一次重新提交进队列。这样一来极端情况下最多延迟几分钟但不会丢。乱序问题在我们场景里影响不大因为订单支付成功通知本质上不依赖严格顺序。但如果你的业务是“先改状态再发短信”这种强顺序链路那就需要在消息体里带一个业务序列号消费端按序列号做排序或者丢弃过期消息而不是单纯依赖队列的有序性。5.3 这次踩坑后总结的几条团队规约这次事故给团队带来的最大收益不是代码优化本身而是几个流程性的改变第一所有涉及队列、线程池、异步处理的代码CR时必须有压测记录或者明确的性能指标预估不能只聊逻辑正确性。第二使用内存队列必须遵守“有界超时降级”三件套不允许在生产代码里出现裸的put()无限阻塞。第三重试逻辑必须和正常消息分离要么用延迟队列要么用专门的待处理表禁止把失败任务放回队列头部。第四条是我个人加上的——任何异步链路上线前必须有队列深度的监控告警和消费延迟的监控告警。没有监控就是睁着眼睛把系统交给运气。5.4 关于“Bug”这件事的一点感想说到“Bug”网上总能看到类似“codex磁盘bug”或者“winsxs bug”这类听着就让人头大的问题但说实话这些离我们日常开发太远了。真正让我们寝食难安的往往是老大写的这段“当时看起来没问题”的队列代码它不报错、不崩溃只是在你最需要它扛住流量的时候安静地堵在那里把所有入口都堵死。排这种Bug最难的从来不是修复而是找到那个让所有线索串起来的核心因果链——队列满了生产者的线程全被卡在入队上回调接口才变慢背后是消费端单条处理吞吐不够。我个人在实际排查和优化过程中最大的体会是遇到性能相关的线上问题永远不要急着改代码。先把线程栈抓下来把压测复现跑出来把数据摆在桌面上再动手。因为只有数据能告诉你真正该优化的是什么——而不是你“感觉”该优化的是什么。这个习惯在代码评审里同样适用看到一段“能用”的代码多问一句“流量翻十倍它还扛得住吗”很多线上事故就根本不会发生。最后再分享一个实用的小技巧排查完类似问题之后把当时的线程dump、压测脚本和优化对比数据留在一个专门的问题追踪文档里。下次再遇到队列或者线程池相关的性能问题先翻这个文档大概率能省下半天排查时间。