6 个下游聚合有 1 个 hang 住,Tomcat 200 个线程全卡死:CompletableFuture 编排的 4 个隐形约定

📅 2026/8/27 16:11:28
6 个下游聚合有 1 个 hang 住,Tomcat 200 个线程全卡死:CompletableFuture 编排的 4 个隐形约定
title: 6 个下游聚合有 1 个 hang 住Tomcat 200 个线程全卡死CompletableFuture 编排的 4 个隐形约定tags: [Java, CompletableFuture, 异步编程, 线程池, 并发]一个下游拖垮整个商品详情页我们的商品详情页是典型的聚合接口一次请求要拿基础信息、库存、价格、促销、评价、店铺信息6 个下游服务。为了压首屏时间两年前就改成了CompletableFuture并行编排P99 从 380ms 降到 140ms当时还写了篇内部分享。问题出在 2026 年 3 月的一次大促预热。评价服务那边做了一次索引重建ES 集群有个节点的磁盘打满查询开始 hang 住——不是报错是不返回。7 分钟后商品详情页整个不可用网关那边的错误率从 0.02% 飙到 89%。奇怪的是评价信息在页面上只是一个「4.8 分 1.2 万条评价」的角标业务上完全可以降级掉不展示。我们代码里明明写了exceptionally。线程 dump 拉下来200 个http-nio-8080-exec-*线程全部停在同一行http-nio-8080-exec-137 #256 daemon waiting on condition java.lang.Thread.State: WAITING (parking) at jdk.internal.misc.Unsafe.park(java.base17.0.9/Native Method) at java.util.concurrent.CompletableFuture$Signaller.block(...) at java.util.concurrent.ForkJoinPool.managedBlock(...) at java.util.concurrent.CompletableFuture.waitingGet(...) at java.util.concurrent.CompletableFuture.join(CompletableFuture.java:2117) at com.xxx.detail.DetailAggregator.aggregate(DetailAggregator.java:64)join()没有超时。这是第一个坑也是最致命的那个。出事的那段编排代码public DetailVO aggregate(long itemId) { CompletableFutureBaseInfo base CompletableFuture .supplyAsync(() - itemClient.getBase(itemId), bizPool); CompletableFutureStock stock CompletableFuture .supplyAsync(() - stockClient.get(itemId), bizPool); CompletableFuturePrice price CompletableFuture .supplyAsync(() - priceClient.get(itemId), bizPool); CompletableFuturePromotion promo CompletableFuture .supplyAsync(() - promoClient.get(itemId), bizPool); CompletableFutureComment comment CompletableFuture .supplyAsync(() - commentClient.get(itemId), bizPool); CompletableFutureShop shop CompletableFuture .supplyAsync(() - shopClient.get(itemId), bizPool); CompletableFuture.allOf(base, stock, price, promo, comment, shop) .exceptionally(e - { // 问题 1挂在 allOf 上 log.warn(聚合部分失败, e); return null; }) .join(); // 问题 2无超时 return build(base.join(), stock.join(), price.join(), promo.join(), comment.join(), shop.join()); // 问题 3 }逐条拆问题 1exceptionally挂在allOf上不等于给每个子任务兜底。allOf返回的是CompletableFutureVoid它只在所有子 future 都完成后才完成。给它挂exceptionally只能捕获「allOf 本身完成时是异常态」这一种情况改变的是 allOf 那条链的结果对comment这个 future 本身的状态毫无影响。所以最后一行comment.join()依然会抛出CompletionException。问题 2join()不接受超时参数。CompletableFuture里带超时的只有get(long, TimeUnit)它抛受检异常TimeoutException而join()是无限等待。ES hang 住意味着commentClient.get()里的 HTTP 调用没有 socket timeout那个任务永远不完成allOf永远不完成join()就永远不返回。Tomcat 线程一个个耗进去200 个耗尽只要几十秒。问题 3最后连着 6 个join()。即使前面的 join 加了超时并成功跳出这 6 个还是会在异常 future 上抛出来。这是典型的「超时机制只做了一半」。源码层面看 join 到底在等什么JDK 17 的CompletableFuture.join()public T join() { Object r; if ((r result) null) r waitingGet(false); // false 表示不可中断 return reportJoin(r); } private Object waitingGet(boolean interruptible) { Signaller q null; boolean queued false; Object r; while ((r result) null) { if (q null) { q new Signaller(interruptible, 0L, 0L); // 注意这两个 0L if (Thread.currentThread() instanceof ForkJoinWorkerThread) ForkJoinPool.helpAsyncBlocker(defaultExecutor(), q); } else if (!queued) queued tryPushStack(q); else { try { ForkJoinPool.managedBlock(q); // 真正 park 的地方 } catch (InterruptedException ie) { /* ... */ } } } ... }new Signaller(interruptible, 0L, 0L)那两个0L分别是nanos和deadline。传 0 意味着Signaller.isReleasable()里那段 deadline 判断整段短路——永远不会因为超时被释放。这就是join()无限等待的实现层面原因。get(timeout, unit)走的是timedGet()那里会算出真实 deadline。顺带说一个很多人忽略的点ForkJoinPool.managedBlock在调用线程是 FJP worker 时会尝试补偿一个新线程避免整个池饿死。但如果调用线程是 Tomcat 线程不是 FJP worker这套补偿机制完全用不上就是硬 park。第四个约定回调在哪个线程跑修复过程中还发现一个隐患。有位同事为了「省一个线程」把促销的后处理写成了这样CompletableFuturePromotion promo CompletableFuture .supplyAsync(() - promoClient.get(itemId), bizPool) .thenApply(p - { // 这里又调了一次 RPC 查券 return couponClient.enrich(p); // 阻塞调用 });thenApply不带 Async的执行线程规则是如果上游 future 在你调thenApply时已经完成回调就在当前线程调用者线程同步执行如果还没完成就在完成它的那个线程执行。前者意味着 Tomcat 线程会直接跑这段阻塞 RPC后者意味着bizPool的线程要跑完 RPC 才能释放。两种情况都不是作者想要的「异步」。三种写法的实际行为对比写法上游已完成时上游未完成时适合场景thenApply(fn)调用者线程执行完成上游的线程执行纯 CPU 的轻量转换如字段映射thenApplyAsync(fn)ForkJoinPool.commonPoolcommonPool短小 CPU 任务不含阻塞thenApplyAsync(fn, pool)指定 pool指定 pool含 IO/阻塞的后处理必须显式指定我的规则很简单回调里只要有任何形式的 IO、锁、sleep一律用thenApplyAsync(fn, 自己的池)。不带 Async 的版本只留给纯内存计算。commonPool 也不能用于阻塞任务——它默认大小是 CPU 核数减一我们的机器是 8 核也就是 7 个线程几个阻塞调用就能占满还会影响到同 JVM 内所有用 commonPool 的地方包括 parallelStream。改完之后的版本public DetailVO aggregate(long itemId) { // 关键每个子任务自带超时 自带降级不依赖外层兜底 CompletableFutureBaseInfo base withFallback( () - itemClient.getBase(itemId), 300, BaseInfo.EMPTY, base); CompletableFutureComment comment withFallback( () - commentClient.get(itemId), 120, Comment.EMPTY, comment); // 其余 4 个同理省略 // allOf 只用来等齐因为每个子 future 都不会异常完成了 CompletableFuture.allOf(base, stock, price, promo, comment, shop) .orTimeout(400, TimeUnit.MILLISECONDS) // 兜底的兜底 .exceptionally(e - null) .join(); return build(base.getNow(BaseInfo.EMPTY), stock.getNow(Stock.EMPTY), price.getNow(Price.EMPTY), promo.getNow(Promotion.EMPTY), comment.getNow(Comment.EMPTY), shop.getNow(Shop.EMPTY)); } private T CompletableFutureT withFallback( SupplierT call, long timeoutMs, T fallback, String tag) { return CompletableFuture.supplyAsync(call, bizPool) .orTimeout(timeoutMs, TimeUnit.MILLISECONDS) // JDK 9 .exceptionally(e - { // 挂在子任务上 metrics.counter(detail.fallback, tag, tag).increment(); log.warn({} 降级, cause{}, tag, e.toString()); return fallback; }); }几个改动点值得单独说orTimeout是 JDK 9 才有的内部用一个单线程的Delayer调度器ScheduledThreadPoolExecutordaemon 线程名CompletableFutureDelayScheduler在到期时把 future 以TimeoutException完成。注意它不会中断正在执行的任务——底层 HTTP 调用还在跑只是结果不要了。所以 socket timeout 该配还得配这是两层防护不是替代关系。exceptionally挂在每个子 future 上返回 fallback 值于是子 future 变成正常完成状态。这样allOf不会异常后面的getNow也拿得到值。最后用getNow(fallback)而不是join()。getNow不阻塞拿不到就用默认值等于给整条链加了最后一道保险。外层还留了orTimeout(400ms)因为 6 个子任务各 300ms 超时理论上最坏也在 300ms 左右完成但如果线程池排队严重任务还没开始执行超时定时器却已经在跑实际耗时可能超预期外层兜一道。复盘数字故障 7 分 12 秒商品详情页错误率峰值 89%估算影响订单约 1.1 万笔。恢复手段是紧急重启 临时摘掉评价服务的注册节点不是代码修复。改造后做了一次故障演练用 iptables DROP 掉评价服务的返回包模拟 hang。改造前 43 秒线程池耗尽改造后接口 P99 从 138ms 变成 262ms等到评价的 120ms 超时但成功率保持 99.97%评价角标显示为默认值。bizPool配置也调了核心 32、最大 64、队列 200、拒绝策略从AbortPolicy改成CallerRunsPolicy之外再包一层降级——直接CallerRuns会把 Tomcat 线程拖进来这点在评审时被指出来了。我的几个取舍判断不要用allOf().join()这种「等齐再取值」的写法组织聚合接口。它把「所有下游都得成功」这个隐含假设写进了代码而聚合接口的本质是「尽力而为」。更合适的模型是每个子任务自带超时和降级主流程只负责等一个总窗口。join()我现在的态度是业务代码里不允许出现裸的join()。团队在 ArchUnit 里加了检查join()调用前必须有orTimeout或completeOnTimeout。这条规则挡住过两次类似写法进主干。CompletableFuture不适合做复杂的分支编排。超过 3 层依赖、带条件分支的场景代码可读性会掉得很快异常传播路径也很难讲清楚。我更建议这种场景用 Reactor 的Mono.zip加onErrorResume或者干脆退回同步 线程池牺牲一点延迟换可维护性。CompletableFuture的甜点区是「扇出几个独立调用然后合并」超出这个范围就该换工具。降级值不能是 null。我们最早的 fallback 返回 null结果build()里到处是 NPE 判断。改成EMPTY常量对象各字段是零值/空串/空集合之后下游渲染逻辑一行判断都不用改。留个问题orTimeout用的是一个全局单线程的Delayer。如果你的服务 QPS 有 5000每个请求注册 6 个orTimeout也就是每秒往那个单线程调度器里塞 3 万个延时任务——你觉得这个调度器会成为瓶颈吗如果会你打算怎么改欢迎在评论区说说你的方案。