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.base@17.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) { CompletableFuture<BaseInfo> base = CompletableFuture .supplyAsync(() -> itemClient.getBase(itemId), bizPool); CompletableFuture<Stock> stock = CompletableFuture .supplyAsync(() -> stockClient.get(itemId), bizPool); CompletableFuture<Price> price = CompletableFuture .supplyAsync(() -> priceClient.get(itemId), bizPool); CompletableFuture<Promotion> promo = CompletableFuture .supplyAsync(() -> promoClient.get(itemId), bizPool); CompletableFuture<Comment> comment = CompletableFuture .supplyAsync(() -> commentClient.get(itemId), bizPool); CompletableFuture<Shop> 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 }逐条拆:
问题 1:exceptionally挂在allOf上,不等于给每个子任务兜底。allOf返回的是CompletableFuture<Void>,它只在所有子 future 都完成后才完成。给它挂exceptionally,只能捕获「allOf 本身完成时是异常态」这一种情况,改变的是 allOf 那条链的结果,对comment这个 future 本身的状态毫无影响。所以最后一行comment.join()依然会抛出CompletionException。
问题 2:join()不接受超时参数。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。
第四个约定:回调在哪个线程跑
修复过程中还发现一个隐患。有位同事为了「省一个线程」,把促销的后处理写成了这样:
CompletableFuture<Promotion> 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.commonPool | commonPool | 短小 CPU 任务,不含阻塞 |
thenApplyAsync(fn, pool) | 指定 pool | 指定 pool | 含 IO/阻塞的后处理,必须显式指定 |
我的规则很简单:回调里只要有任何形式的 IO、锁、sleep,一律用thenApplyAsync(fn, 自己的池)。不带 Async 的版本只留给纯内存计算。commonPool 也不能用于阻塞任务——它默认大小是 CPU 核数减一,我们的机器是 8 核,也就是 7 个线程,几个阻塞调用就能占满,还会影响到同 JVM 内所有用 commonPool 的地方(包括 parallelStream)。
改完之后的版本
public DetailVO aggregate(long itemId) { // 关键:每个子任务自带超时 + 自带降级,不依赖外层兜底 CompletableFuture<BaseInfo> base = withFallback( () -> itemClient.getBase(itemId), 300, BaseInfo.EMPTY, "base"); CompletableFuture<Comment> 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> CompletableFuture<T> withFallback( Supplier<T> 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调度器(ScheduledThreadPoolExecutor,daemon 线程名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 万个延时任务——你觉得这个调度器会成为瓶颈吗?如果会,你打算怎么改?欢迎在评论区说说你的方案。