news 2026/8/27 16:11:26

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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.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分别是nanosdeadline。传 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.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) { // 关键:每个子任务自带超时 + 自带降级,不依赖外层兜底 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()调用前必须有orTimeoutcompleteOnTimeout。这条规则挡住过两次类似写法进主干。

CompletableFuture不适合做复杂的分支编排。超过 3 层依赖、带条件分支的场景,代码可读性会掉得很快,异常传播路径也很难讲清楚。我更建议这种场景用 Reactor 的Mono.ziponErrorResume,或者干脆退回同步 + 线程池,牺牲一点延迟换可维护性。CompletableFuture的甜点区是「扇出几个独立调用,然后合并」,超出这个范围就该换工具。

降级值不能是 null。我们最早的 fallback 返回 null,结果build()里到处是 NPE 判断。改成EMPTY常量对象(各字段是零值/空串/空集合)之后,下游渲染逻辑一行判断都不用改。

留个问题

orTimeout用的是一个全局单线程的Delayer。如果你的服务 QPS 有 5000,每个请求注册 6 个orTimeout,也就是每秒往那个单线程调度器里塞 3 万个延时任务——你觉得这个调度器会成为瓶颈吗?如果会,你打算怎么改?欢迎在评论区说说你的方案。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/27 16:06:33

一条链接搞定B站视频下载与AI总结

一条链接搞定B站视频下载与AI总结 【免费下载链接】BiliTools 本项目已停止维护。 项目地址: https://gitcode.com/GitHub_Trending/bilit/BiliTools 你看过两小时的网课录像&#xff0c;想回看某个片段时找不到位置&#xff0c;只能从头再放一遍。BiliTools 是一款跨平…

作者头像 李华
网站建设 2026/8/27 16:05:17

低剖面180W AC-DC电源设计:从效率到散热的全流程解析

这两年做系统集成的工程师应该都有同感&#xff1a;设备越做越薄&#xff0c;留给电源的高度空间越来越紧。机箱从3U压到1U&#xff0c;LED显示箱体从80mm压到60mm&#xff0c;控制柜里导轨排得密密麻麻。在这种大趋势下&#xff0c;180W级别的AC-DC电源能把高度压进20mm出头&a…

作者头像 李华
网站建设 2026/8/27 16:03:31

MSLab 入门指南:用 3 条 PowerShell 命令搭出 Azure Local 测试集群

MSLab 入门指南&#xff1a;用 3 条 PowerShell 命令搭出 Azure Local 测试集群 【免费下载链接】MSLab Azure Local (formerly Azure Stack HCI), Windows 10 and Windows Server rapid lab deployment scripts 项目地址: https://gitcode.com/gh_mirrors/ms/MSLab MSL…

作者头像 李华