第一次在生产环境里用 CompletableFuture 把五个串行接口压成一次并行调用,接口 P95 从 480ms 掉到 130ms,那种感觉确实挺爽。但同样也是它,让我在某个凌晨两点对着一个永远不返回的 join() 干瞪眼,最后发现是自定义线程池的队列满了、拒绝策略还是 CallerRunsPolicy,主线程被自己提交的任务给堵死了。CompletableFuture 这个工具就是这样,上手半小时就能写出能跑的代码,但要写出在生产环境里扛得住、查得动、改得动的异步编排,中间隔着好几层坑。这篇东西我打算把它讲透:从它到底补上了 Future 的哪些短板,到 thenApply、thenCompose、thenCombine 这几个最容易混的 API 该怎么选,再到线程池参数怎么算、超时怎么设、异常为什么老是"凭空消失"、线上线程池打满该怎么排查。文中的代码片段都是可以直接拿去改的骨架,参数计算过程我会一步步写出来,不看结论也能自己推。适合已经会写 Java、但对异步编排还停留在"抄一段能用就行"阶段的同学,也适合那些代码已经上了生产、正在被莫名其妙的超时和线程堆积折磨的同学。
1. 为什么值得把 CompletableFuture 当成主力异步编排工具
1.1 从 Future 那点憋屈说起
Java 5 就把 Future 塞进标准库了,初衷很清楚:把耗时操作丢给别的线程,主线程该干嘛干嘛。但真正在业务代码里用过的人心里都有数,Future 的 API 少得可怜——提交、取消、判断是否完成、get 拿结果,就这四件事。它的定位其实是一个"结果占位符",压根就不是编排工具。这个定位差异,直接导致了它在稍微复杂一点的场景下就力不从心。
举个具体的例子。你在做一个商品详情页,需要查商品基础信息、库存、价格、评价数量、推荐列表,这五份数据来自五个不同的下游服务。用 Future 你只能 submit 五次,然后五次 get。这个 get 是阻塞的,谁先完成谁后完成你控制不了,只能按固定的顺序一个个等。五个接口平均各 80ms,最慢的 200ms,总体耗时就是五者之和再叠上网络抖动,串行的本质没变,你只是把等待挪到了不同的 get 调用上而已。
更要命的是依赖关系。假设推荐列表的请求参数里需要带上商品类目,而类目要从基础信息里取,这就变成了两段式调用:先拿基础信息,再拿推荐。用 Future 你只能等第一个 get 返回,拿到类目,再 submit 第二个任务,再加到列表里等。代码写成什么样子可以想象,一堆中间变量、一堆 try-catch、一堆手动维护的 List ,可读性基本为零。
异常处理也是老问题。Future.get() 抛出的是 ExecutionException,真正的业务异常被包了一层,你得 catch 住之后再调用 getCause() 才能拿到原始异常,如果中间还有别的框架又包了一层,那就得剥好几层。日志里打出来的栈信息经常是一堆 java.util.concurrent.ExecutionException,翻半天看不到真正的报错点。
还有一个很隐蔽的限制:Future 没法手动完成。有些场景下结果不是由计算产生的,而是由外部事件触发的,比如一个回调、一条消息、一次信号通知。这种场景你用 Future 就只能开一个线程去阻塞等待,白白占着一个线程不放,线程池稍微小一点就直接被占满了。这就是典型的"用阻塞的方式实现异步",听起来荒谬,但在 Future 时代确实是常规操作。
1.2 CompletableFuture 补上的三块拼图
CompletableFuture 从 Java 8 进入标准库,它同时实现了 Future 和 CompletionStage 两个接口。这两个接口恰好对应了两层用途:Future 那一层负责"最终拿到结果",CompletionStage 那一层负责"把多个阶段串起来、拼起来、合起来"。理解了这层双身份,很多 API 的设计意图就通了。
它补上的第一块拼图是回调式编程。thenApply、thenAccept、thenRun 这一族方法的作用是:在当前阶段完成的那一瞬间,自动触发后续动作,完全不需要有线程在那里阻塞等待。这一点带来的直接好处是线程利用率。一个线程提交完任务就可以立刻回去处理别的请求,等结果就绪时由完成该任务的线程顺手把回调执行掉,或者异步交给另一个线程池执行。整个过程没有线程被空转占用。
第二块拼图是组合能力,也是它真正被称为"编排"工具的原因。thenCompose 用来串联两个有前后依赖的异步操作,把 CompletableFuture<CompletableFuture > 这种嵌套结构扁平化;thenCombine 用来把两个彼此独立的异步结果合并成一个;allOf 和 anyOf 用来聚合一批任务。有了这几个方法,前面那个"先查基础信息再查推荐"的两段式调用,就能写成一个连贯的链式表达式。第三块拼图是手动完成,complete() 和 completeExceptionally() 这两个方法让 CompletableFuture 成了一个可以被外部驱动的状态容器,超时控制、事件驱动、回调桥接这些场景都靠它实现。
1.3 什么场景适合用,什么场景别硬上
我不建议把 CompletableFuture 当成银弹到处套。异步带来的复杂度是实打实的:线程上下文会丢、异常栈会变得又长又难读、调试的时候断点跳来跳去、线程池管理不当会直接拖垮整个应用。一个只有两个串行调用、总耗时 120ms 的接口,改成并行之后变成 90ms,省下来的 30ms 根本抵不上你为它引入的维护成本。
我的判断标准大概是这样。适合用的场景有三个特征:调用的下游数量多(三个以上)、彼此之间没有强依赖(或者依赖关系是清晰的树状)、单个调用的耗时有明显波动(说明有 I/O 等待可以榨取)。典型的例子就是详情页聚合、订单列表批量补数据、报表多维度统计、风控规则并行执行。这些场景下并行化带来的收益是成倍的,代码复杂度换来的性能提升划算。
不适合的场景也很明确。CPU 密集型任务本来就该用并行流或者专门的线程池,套 CompletableFuture 只会让代码更绕。完全串行依赖的流程,比如"扣款成功才能发货",异步化毫无意义,反而把事务边界搞乱了。还有一种容易被忽略的情况是低频接口,比如一天调用几百次的运营后台功能,为了它去折腾异步编排,收益基本可以忽略,写清楚比写快更重要。
2. 核心 API 全景拆解:分清每个方法到底在干什么
2.1 创建任务:supplyAsync 和 runAsync 的差别不只是有没有返回值
CompletableFuture 的入口方法主要就四个重载:runAsync(Runnable)、runAsync(Runnable, Executor)、supplyAsync(Supplier)、supplyAsync(Supplier, Executor)。前两个没有返回值,后两个有返回值,这个差别谁都看得出来。真正值得说的是另一个差别:不传 Executor 的时候,任务会跑到 ForkJoinPool.commonPool() 里,而 commonPool 的默认线程数是 CPU 核数减一。
这意味着在一台四核的机器上,你的所有没指定线程池的 CompletableFuture 任务,都挤在三个线程里抢。一旦链路上有任何一个环节发生了阻塞(比如一次没设超时的 HTTP 调用),这三个线程很快就被占满,然后所有依赖 commonPool 的代码一起卡住。这不是理论风险,我在项目里见过不止一次:某个同事写了个工具类,里面用 supplyAsync 没传线程池,量一上来整个应用的异步链路全堵住,日志里看不到任何报错,就是响应越来越慢。
所以我给自己定的规矩很硬:业务代码里的 supplyAsync 和 runAsync,Executor 参数一律不能省。哪怕只是临时用一下,也要显式传一个自己管理的线程池,哪怕就是 Executors.newFixedThreadPool(4) 这种简单实现,至少它和 commonPool 是隔离的,出问题的时候边界清楚。
// 不推荐:默认跑在 commonPool,线程数受限且与全应用共享 CompletableFuture<String> bad = CompletableFuture.supplyAsync(() -> queryRemote()); // 推荐:显式指定业务隔离的线程池 Executor bizPool = new ThreadPoolExecutor( 8, 16, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(256), new NamedThreadFactory("detail-fetch-"), new ThreadPoolExecutor.CallerRunsPolicy()); CompletableFuture<String> good = CompletableFuture.supplyAsync(() -> queryRemote(), bizPool);提示:commonPool 并不是完全不能用。纯 CPU 计算、没有阻塞、任务粒度小且数量可控的场景,用它其实是合适的。关键在于你要清楚自己在用哪个池子,而不是"忘了传"。
2.2 结果转换:thenApply、thenAccept、thenRun 三个门派的边界
这三兄弟是链式写法里出现频率最高的,但它们的分工完全不同,混用会导致编译不过或者行为不符合预期。判断方法很简单,看输入和输出。
thenApply(Function<T,R>):拿到上一步的结果,处理之后返回一个新结果。它是"转换",会改变阶段的值。适合做数据加工,比如把远程返回的 JSON 字符串转成对象、把金额单位从分转成元。
thenAccept(Consumer ):拿到上一步的结果,消费掉,但不返回值。它的返回类型是 CompletableFuture 。
thenRun(Runnable):连上一步的结果都不需要,只是在前面完成后触发一个动作,同样返回 CompletableFuture 。
这三者还有一个共同的后缀版本:thenApplyAsync、thenAcceptAsync、thenRunAsync。带 Async 的方法会把后续动作交给线程池执行,不带 Async 的方法则在"完成前一个阶段的那个线程"上直接执行。这一点非常关键,也是最容易踩坑的地方。
假设你的链路是这样:supplyAsync(fetch, poolA).thenApply(this::parse)。如果 fetch 是在 poolA 的线程 X 上完成的,那么 parse 也会在 X 上执行。如果 parse 是个耗时 50ms 的 JSON 解析,那 X 就被占用了 50ms 不能干别的。更极端的情况是,如果你的 parse 里有阻塞调用,那就是在浪费线程池的宝贵线程。所以我的经验是:轻量级、纯内存的转换用不带 Async 的版本(省一次线程切换,性能反而更好),任何可能阻塞或者耗时不确定的操作,一律用带 Async 的版本,并且显式指定线程池。
CompletableFuture<Order> future = CompletableFuture .supplyAsync(() -> fetchOrderRaw(orderId), ioPool) // I/O 在 ioPool .thenApply(this::parseOrder) // 内存解析,线程复用即可 .thenApplyAsync(this::enrichWithPromotion, cpuPool) // 涉及计算,切到 cpuPool .thenApply(order -> { log.info("order assembled: {}", order.getId()); return order; });2.3 thenCompose 与 thenCombine:编排里最容易混的一对
这两个名字长得像、签名也像,但语义差得远,几乎是面试和 code review 里的必考题。
thenCompose(Function<T, CompletionStage>):处理的是"依赖关系"。当你的下一步操作需要用到上一步的结果作为参数,而且这个下一步本身又是异步的,就用它。它的作用是扁平化,把 CompletableFuture<CompletableFuture> 压成 CompletableFuture。
thenCombine(CompletionStageother, BiFunction<T,U,V>):处理的是"并列关系"。当你有两个互相独立的异步任务,需要把两个结果合起来做一件事,就用它。两个任务会同时跑,谁先谁后不重要。
用一个具体的例子区分。查订单详情的时候,"从订单服务拿订单基本信息"和"从用户服务拿用户资料"这两个动作互不依赖,可以并行,最后一起组装成页面数据,这就是 thenCombine 的场景。而"先拿订单基本信息,从里面取出用户 ID,再拿用户资料",这就是 thenCompose 的场景,因为第二步的参数依赖第一步的产出。
如果该用 thenCompose 的地方误用了 thenApply,你会得到 CompletableFuture<CompletableFuture > 这种嵌套类型,写起来很别扭,而且外层 future 完成的时候内层任务可能还没跑完,很容易出现"拿到的是个没完成的 Future"这种诡异 bug。反过来,如果该用 thenCombine 的地方用成了顺序的 thenCompose,功能上不会错,但两个本该并行的调用被强制串行了,性能白白损失一半。
// 场景一:并行 + 合并,用 thenCombine CompletableFuture<Order> orderF = CompletableFuture.supplyAsync(() -> orderService.get(orderId), ioPool); CompletableFuture<User> userF = CompletableFuture.supplyAsync(() -> userService.get(userId), ioPool); CompletableFuture<DetailVO> vo = orderF.thenCombine(userF, this::assembleDetail); // 场景二:依赖 + 串联,用 thenCompose CompletableFuture<User> userF2 = CompletableFuture .supplyAsync(() -> orderService.get(orderId), ioPool) .thenCompose(order -> CompletableFuture.supplyAsync( () -> userService.get(order.getUserId()), ioPool));2.4 allOf / anyOf:多任务聚合的真实行为
当并行任务从两个变成五个、十个,thenCombine 就不好写了,因为它只接受两个参数。这时候要用 allOf。
CompletableFuture.allOf(f1, f2, f3, ...) 返回一个 CompletableFuture ,它在所有传入的任务都完成后才完成。注意它返回的是 Void,也就是说它只告诉你"都完事了",但不携带任何结果数据。这是初学者最容易懵的地方——很多人以为 allOf 会返回一个结果列表,结果发现类型是 Void,不知道怎么取值。
正确的用法是:先保存每个子任务的引用,用 allOf 等待全部完成,然后从各个引用上逐个 join 取结果。
CompletableFuture<BaseInfo> base = CompletableFuture.supplyAsync(() -> baseService.get(sku), ioPool); CompletableFuture<Stock> stock = CompletableFuture.supplyAsync(() -> stockService.get(sku), ioPool); CompletableFuture<Price> price = CompletableFuture.supplyAsync(() -> priceService.get(sku), ioPool); CompletableFuture<Comments> cmt = CompletableFuture.supplyAsync(() -> cmtService.get(sku), ioPool); CompletableFuture<Void> all = CompletableFuture.allOf(base, stock, price, cmt); CompletableFuture<DetailVO> result = all.thenApply(v -> { DetailVO vo = new DetailVO(); vo.setBase(base.join()); vo.setStock(stock.join()); vo.setPrice(price.join()); vo.setComments(cmt.join()); return vo; });这里每个 join() 都不会真的阻塞,因为 allOf 已经保证了所有任务都已完成,join 只是把结果取出来。这个细节很重要,很多人担心 join 会阻塞,其实在 allOf 的回调里它是立即返回的。
注意:如果其中任何一个子任务抛出了异常,allOf 返回的 future 会立刻以异常完成,不会等剩下的任务。而剩下的任务其实还在跑,只是没人管了。这个行为在设计降级策略的时候要考虑到。
anyOf 的行为是"任意一个完成就返回",它返回的是 CompletableFuture
还有一个聚合的替代方案值得提一下:如果你用的是 Java 9 及以上,可以试试 CompletableFuture 的延迟组合,或者干脆用 Stream API 配合自定义的 collector 做批量聚合。不过在 Java 8 环境里,allOf 加 join 这个组合依然是最稳、最通用的写法。
3. 线程池选型与参数计算实操
3.1 默认的 commonPool 为什么不能当生产主力
前面简单提过,这里展开说。ForkJoinPool.commonPool() 是 JVM 全局共享的,默认并行度是 Runtime.getRuntime().availableProcessors() - 1。在容器化环境中,这个值取决于 JVM 是否感知到了容器配额。Java 8u191 之前,JVM 读的是宿主机的核数而不是容器限额,一个限制在两核的 Pod 可能看到六十四核,commonPool 就开了六十三个线程,上下文切换开销直接把性能拖垮。即使版本够新,commonPool 仍然有三个硬伤。
第一,所有模块共享。你引入的第三方库、框架内部的异步实现、自己写的业务代码,都用同一个池子。任何一个地方发生阻塞,都会影响到其他所有使用者。这种"故障串联"在生产环境里是非常危险的。
第二,没有队列控制。commonPool 内部的任务队列是无界的,任务提交速度超过消费速度时,任务会无限堆积,直到内存被撑爆。你没法给它设一个上限来触发快速失败。
第三,不可监控。你拿不到它的活跃线程数、队列长度、拒绝次数,出了问题只能靠 jstack 抓现场,排查效率极低。
所以结论很直接:生产环境的业务异步,一律用自建线程池,把线程数、队列长度、拒绝策略、线程名前缀、监控埋点全部握在自己手里。
3.2 线程数到底开多少:从公式推导到场景修正
线程池大小的经典公式是这样的:
线程数 = CPU 核数 × 目标 CPU 利用率 × (1 + 等待时间 / 计算时间)这个公式的来源是把线程分成两种状态:跑在 CPU 上的时间和等待时间(I/O、锁、网络)。等待占比越高,需要越多的线程来把 CPU 喂饱。举个具体的计算过程。
假设部署在四核机器上,目标 CPU 利用率定 0.8,避免跑满导致排队。一次下游调用中,序列化请求、拼装参数、解析响应的纯计算部分加起来约 5ms,而网络往返加下游处理约 95ms。那么等待时间与计算时间的比值就是 95 / 5 = 19。
代入公式:4 × 0.8 × (1 + 19) = 64。也就是说理论上需要六十四个线程才能让 CPU 保持 80% 的利用率。这个数字看起来很大,但在纯 I/O 等待的场景下是合理的——线程大部分时间都在等网络,CPU 是空闲的。
不过公式给的是理论上限,实际落地的数值还要做几轮修正。
第一轮,看下游的承载能力。六十四个线程意味着对下游可能有六十四个并发请求,如果下游是个承载能力只有三十并发的老服务,你开六十四就是把它打死,然后大家一起超时。这时候必须按下游容量倒推,宁可让上游排队。
第二轮,看总资源。每个线程大约占 1MB 栈空间(默认 -Xss),六十四线程就是 64MB,加上堆内存,得算一下容器的内存限制够不够。另外还要考虑你的服务同时对外提供多少并发,接受的请求数乘以每个请求的线程占用,不能超过池子容量太多,否则大量请求会在队列里排队,延迟反而更高。
第三轮,做压测确认。公式算出来的永远只是起点,一定要在生产同规格的环境上压一轮,看几个关键指标:线程池活跃数、队列长度、任务平均等待时间、下游错误率。逐步调整线程数,找到"吞吐量不再上升、但延迟开始恶化"的那个拐点,往回收 10% 到 20% 就是比较合适的配置。
我实际项目中用得比较多的一个经验值是:纯 I/O 型的远程调用线程池,线程数取 CPU 核数的 8 到 16 倍,队列长度取线程数的 2 到 4 倍,剩下的用拒绝策略兜底。这个范围不是拍脑袋,是在几个不同量级的服务上反复压测收敛出来的。
3.3 线程池隔离与监控埋点
线程池隔离的原则很简单:不同性质的任务不能共用一个池子。我通常会按这个维度切分:远程调用类(HTTP、RPC)一个池子,本地缓存访问如果用了阻塞客户端再开一个小池子,CPU 密集的计算一个池子,定时任务一个池子。这样做的价值在于故障隔离——某一个下游变慢,只会打满它对应的那个池子,不会影响其他链路。
线程命名前缀是低成本的排查利器。给每个池子的线程起一个有意义的前缀,比如 "detail-rpc-"、"report-calc-",线上 dump 线程栈的时候一眼就能看出是谁在忙、谁堵住了。这个习惯建议从项目第一天就养成,后期补非常痛苦。
监控埋点至少要覆盖这几个指标,用 Micrometer 或你项目里现成的监控库都能做:
| 指标 | 采集方式 | 告警阈值建议 |
|---|---|---|
| 活跃线程数 | ThreadPoolExecutor.getActiveCount() | 持续超过核心线程数的 80% |
| 队列长度 | getQueue().size() | 超过队列容量的 60% |
| 已完成任务数 | getCompletedTaskCount() | 增速骤降需要关注 |
| 拒绝任务数 | 自定义 RejectedExecutionHandler 计数 | 大于 0 即告警 |
| 任务等待时间 | 提交时打时间戳,执行时算差值 | P99 超过 100ms |
自定义拒绝策略的地方要多说一句。ThreadPoolExecutor 内置的四种策略里,AbortPolicy 会抛 RejectedExecutionException,CallerRunsPolicy 会让提交任务的线程自己执行,也就是把压力传导回上游。这两者在异步编排里的表现完全不同。用 AbortPolicy,异常会沿着 CompletableFuture 的链路传播,最后被 exceptionally 捕获,你能明确知道发生了拒绝。用 CallerRunsPolicy,任务被主线程执行,看似"更平滑",但如果主线程本身就是 Tomcat 的工作线程,那就等于把 Web 容器的线程也拖进了这个池子的逻辑里,容易引发连锁的线程耗尽。
我个人的选择是:核心业务链路用 AbortPolicy 加明确的降级逻辑,非核心的、可以容忍延迟的用 CallerRunsPolicy。同时在拒绝处理器里打一条带线程池名字和当前状态的日志,这是事后定位问题的关键证据。
4. 异步链路上的异常与超时
4.1 异常在链路上是怎么"走"的
CompletableFuture 的异常传播机制和同步代码完全不一样,这也是它最容易让人困惑的地方。核心规则是:一旦某个阶段以异常完成,后续所有不带恢复能力的阶段都会直接跳过,异常一路向后传播,直到遇到一个能处理它的方法。
具体来说,如果链路是 supplyAsync(...).thenApply(...).thenAccept(...),而 supplyAsync 里抛了异常,那么 thenApply 和 thenAccept 都不会执行,最终这个 CompletableFuture 是以异常状态结束的。这时候如果你调用 join(),会抛 CompletionException;调用 get(),会抛 ExecutionException,两者的原始异常都藏在 cause 里。
这里有一个非常致命的坑:如果你不调用 join 或 get,也没有在任何地方注册异常处理,那么这个异常就彻底消失了。没有任何日志,没有任何告警,代码静悄悄地什么都不做。这是异步编程里最难查的一类问题——业务方说"这个功能没生效",你看代码逻辑完全正确,最后发现是某一步抛了异常没人管。
所以我的硬性要求是:每一条 CompletableFuture 链路的末端,必须有明确的异常处理,而且异常处理里必须有日志。哪怕只是 .exceptionally(ex -> { log.error("xxx failed", ex); return null; }) 这么简单,也比什么都没有强得多。
4.2 exceptionally、handle、whenComplete 三个方法怎么选
这三个方法都能接触异常,但用途不同,选错了要么编译不过,要么行为不符合预期。
exceptionally(Function<Throwable, T>):只在发生异常时触发,返回一个同类型的兜底值,把异常"吞掉",让链路恢复成正常完成状态。适合做降级,比如远程调用失败时返回一个默认的空对象。
handle(BiFunction<T, Throwable, R>):无论成功还是失败都会触发,接收结果和异常两个参数,返回一个新的值。它既能处理异常也能处理正常结果,而且可以改变返回类型。适合"不管成不成功,我都要重新组织一下返回结构"的场景。
whenComplete(BiConsumer<T, Throwable>):无论成功还是失败都会触发,但不改变结果。它像是 finally 块,用来做收尾动作,比如释放资源、记录耗时、上报监控。它返回的 CompletableFuture 仍然保持原来的结果或异常状态,异常会继续往后传。
一个常见的误用是:用 exceptionally 做了降级,然后以为链路后面的 whenComplete 还能感知到原始异常——实际上感知不到了,因为异常已经被前面吞掉了。所以在设计恢复顺序的时候,恢复操作要放在收尾操作之后,或者干脆用 handle 一次性把两者都处理掉。
CompletableFuture<DetailVO> safe = rawFuture .exceptionally(ex -> { log.error("detail assemble failed, orderId={}", orderId, ex); metrics.counter("detail.fail").increment(); return DetailVO.empty(); // 降级兜底 }) .whenComplete((vo, ex) -> { // 走到这里时 ex 一定是 null,因为上面已经恢复过了 log.info("detail assembled, cost={}ms", System.currentTimeMillis() - start); });4.3 超时控制:orTimeout 与 completeOnTimeout 的差别
没有超时的异步链路等于没有链路。一个卡住的下游会顺着链路一路把线程池占满,然后整个服务雪崩。CompletableFuture 在 Java 9 之后提供了两个原生超时方法,差别值得说清楚。
orTimeout(long timeout, TimeUnit unit):到达超时时间还没完成,就让这个 future 以 TimeoutException 异常完成。链路会因为异常往下走,最终触发你的 exceptionally 降级逻辑。
completeOnTimeout(T value, long timeout, TimeUnit unit):到达超时时间还没完成,就用给定的默认值让 future 正常完成。它不会产生异常,链路继续往下走,只是拿到的值是兜底值。
选择依据很简单:如果你的业务允许"超时了就用默认值顶上",用 completeOnTimeout 会让代码更干净;如果需要明确区分"正常返回"和"超时降级"这两种情况,用 orTimeout 配合 exceptionally。
// Java 9+ CompletableFuture<Price> priceF = CompletableFuture .supplyAsync(() -> priceService.get(sku), ioPool) .orTimeout(200, TimeUnit.MILLISECONDS) .exceptionally(ex -> { log.warn("price timeout, fallback to default, sku={}", sku); return Price.defaultValue(); });如果你还在 Java 8,没有这两个方法,得自己实现。标准做法是用一个定时器在超时后调用 completeExceptionally,或者用 ScheduledExecutorService 配合一个中间 CompletableFuture。这个自己实现的版本要注意清理定时任务,否则大量超时的任务会堆积成内存泄漏。
注意:orTimeout 只是让 CompletableFuture 对象进入完成状态,它并不会中断底层正在执行的任务。下面那个还在阻塞的 HTTP 调用依然会继续跑,直到它自己超时或者返回。所以真正要控制资源占用,还是得给底层的 HTTP 客户端、RPC 框架设置连接超时和读取超时,双重保险才稳。
5. 实战:订单详情页并行聚合改造
5.1 改造前的耗时基线
这个案例来自一个真实的订单详情接口。改造前它是一个典型的串行实现:先查订单主表拿基本信息(约 60ms),再根据订单里的用户 ID 查用户信息(约 50ms),再查商品信息(约 70ms),再查物流状态(约 120ms),最后查可用的营销活动(约 90ms)。五个调用串起来,加上自身的处理逻辑,P50 大约 400ms,P95 到了 520ms。
问题在于这五个调用里,物流和营销这两个是比较慢的下游,而且调用量一大的时候波动很明显。它们在链路的末端,前面的调用把最坏延迟都累积进去了,导致尾延迟特别难看。这是串行聚合的典型症状:平均值还行,P99 一塌糊涂。
优化的空间也很清楚:这五个调用里,用户信息和商品信息只依赖订单里的 ID,和订单主表的查询结果没有强依赖关系,完全可以并行。物流和营销也是同理。真正有依赖的只有"先查订单拿到 ID"这一步,其他四步都能并行。
5.2 依赖关系拆解与编排设计
拆解之后,整个流程分成两段。
第一段是订单主表查询,必须最先执行,因为后面的用户 ID、商品 ID 都从它里面取。这一段没法并行。
第二段是四个可以并行的查询:用户信息、商品信息、物流状态、营销活动。这四个都拿到订单 ID 之后就可以同时发起。
第二段之后还有一个组装阶段,需要四个结果都到齐才能拼装 DetailVO。这一步必须等全部完成。
对应的编排结构就是:先 supplyAsync 查订单,然后 thenCompose 进入并行阶段;并行阶段用四个 supplyAsync 发起四个任务,用 allOf 等待全部完成,然后 thenApply 组装。
还有一个额外的设计考虑:物流和营销这两个慢且非核心的调用,应该做超时降级。物流查不到就显示"暂无物流信息",营销查不到就不展示营销模块,不能因为它们拖慢整个页面。用户信息和商品信息是核心,超时的话应该让整个请求失败,返回错误提示。
5.3 代码落地:聚合、降级、超时
public DetailVO getOrderDetail(String orderId) { long start = System.currentTimeMillis(); CompletableFuture<DetailVO> chain = CompletableFuture // 第一段:查订单主表,后续所有调用都依赖它 .supplyAsync(() -> orderService.getOrder(orderId), rpcPool) // 进入并行阶段 .thenCompose(order -> { String userId = order.getUserId(); String skuId = order.getSkuId(); // 核心数据:超时则整体失败 CompletableFuture<User> userF = CompletableFuture .supplyAsync(() -> userService.get(userId), rpcPool) .orTimeout(300, TimeUnit.MILLISECONDS); CompletableFuture<Item> itemF = CompletableFuture .supplyAsync(() -> itemService.get(skuId), rpcPool) .orTimeout(300, TimeUnit.MILLISECONDS); // 非核心数据:超时降级为空,不影响整体 CompletableFuture<Logistics> logisticsF = CompletableFuture .supplyAsync(() -> logisticsService.query(orderId), rpcPool) .orTimeout(250, TimeUnit.MILLISECONDS) .exceptionally(ex -> { log.warn("logistics timeout, degrade, orderId={}", orderId); return Logistics.empty(); }); CompletableFuture<Promotion> promoF = CompletableFuture .supplyAsync(() -> promoService.query(userId, skuId), rpcPool) .orTimeout(250, TimeUnit.MILLISECONDS) .exceptionally(ex -> { log.warn("promotion timeout, degrade, orderId={}", orderId); return Promotion.empty(); }); // 等四个都完成再组装 return CompletableFuture .allOf(userF, itemF, logisticsF, promoF) .thenApply(v -> { DetailVO vo = new DetailVO(); vo.setOrder(order); vo.setUser(userF.join()); vo.setItem(itemF.join()); vo.setLogistics(logisticsF.join()); vo.setPromotion(promoF.join()); return vo; }); }) .exceptionally(ex -> { log.error("order detail failed, orderId={}", orderId, ex); throw new BizException("订单详情加载失败,请稍后重试"); }) .whenComplete((vo, ex) -> { long cost = System.currentTimeMillis() - start; metrics.timer("order.detail.cost").record(cost, TimeUnit.MILLISECONDS); log.info("order detail done, orderId={}, cost={}ms", orderId, cost); }); return chain.join(); }有几处细节值得单独说明。第一,rpcPool 是我专门为这个接口配的线程池,核心线程数 16,最大 32,队列 128,线程名前缀 "order-detail-"。这个数字是根据前面的公式和几轮压测收敛出来的,不是随便填的。
第二,超时时间的选择也有依据。物流和营销的日常 P99 大约是 180ms 和 150ms,我给的 250ms 留了约 40% 的余量,既能覆盖大部分正常情况,又能在下游劣化的时候及时切走。而核心数据的 300ms 是基于订单主表 P99 约 200ms 估算的,留出更宽一点的空间。
第三,最外层的 exceptionally 重新抛了一个业务异常。这一点很重要——如果不抛,接口会返回一个空的 DetailVO,前端会展示成一个空白页,用户完全不知道发生了什么。抛出明确异常让上层统一处理成友好的提示,是更好的体验。
5.4 压测结果对比
在同规格的预发环境,用同样的流量模型(500 QPS,持续 10 分钟)压测,结果如下:
| 指标 | 串行版本 | 并行版本 | 变化 |
|---|---|---|---|
| P50 耗时 | 402ms | 128ms | -68% |
| P95 耗时 | 520ms | 176ms | -66% |
| P99 耗时 | 890ms | 310ms | -65% |
| 平均 CPU 使用率 | 22% | 26% | +4pt |
| 线程池活跃数峰值 | 8 | 21 | — |
| 下游错误率 | 0.3% | 0.4% | 基本持平 |
P99 从 890ms 降到 310ms 是收益最大的一块,因为它把原来串行累加的最坏延迟给消掉了。CPU 只涨了四个百分点,因为这个接口绝大部分时间都在等 I/O。下游错误率基本没变,说明并行并没有给下游增加过大的压力——这一点很关键,如果并行之后错误率明显上升,就得回头检查线程数和下游容量的匹配关系。
还有一个不在表格里但很明显的收益:代码结构变得清晰了。串行版本里那五个调用各自的 try-catch、各自的判空、各自的默认值处理,散落在几十行代码里。改成编排之后,每个调用的失败策略都和它自己写在一起,读代码的时候不需要上下跳。
6. 常见问题与排查技巧实录
6.1 常见问题速查表
| 现象 | 大概率原因 | 排查手段 | 解决方向 |
|---|---|---|---|
| 链路完全不执行,无日志无异常 | 某阶段抛异常且无处理 | 在链路末端加 whenComplete 打日志 | 补 exceptionally,每段链路必须有异常出口 |
| 接口偶发超时,但下游监控正常 | 使用了 commonPool 被其他任务占满 | 检查是否漏传 Executor 参数 | 所有异步任务显式指定业务线程池 |
| 线程池活跃数长期打满 | 下游变慢或有阻塞调用 | 打线程栈看线程状态 | 设超时 + 降级 + 按下游容量调线程数 |
| join() 长时间不返回 | 前置任务被阻塞或死锁 | jstack 找线程栈,看谁在等谁 | 加超时,检查是否有嵌套等待 |
| 任务结果丢失、数据不完整 | allOf 里有任务异常后其他任务被忽略 | 检查每个子任务的异常处理 | 子任务单独做降级,别让一个失败影响整体 |
| 日志里 TraceId 是空的 | ThreadLocal 没有跨线程传递 | 检查 MDC、TraceContext 的传递 | 用装饰器包装线程池或手动透传上下文 |
6.2 线程上下文丢失:一个几乎人人都会踩的坑
ThreadLocal 是和线程绑定的,这个大家都知道。但很多人写异步代码的时候会下意识忘记:主线程里 set 进去的值,新线程里读不到。
最典型的表现是日志里的 TraceId 变成空的,链路追踪断掉,出问题的时候根本串不起一次请求的完整调用链。另一类表现是用户上下文丢失,异步任务里拿不到当前登录用户,权限校验直接失败。
解决思路有两种。一种是手动透传:在提交任务之前把需要的值从 ThreadLocal 里取出来,作为方法参数传到异步任务里,在任务开头重新 set 进去。这种方式最直白,但代码里会到处都是这种透传样板,维护起来烦。
另一种是包装线程池。写一个 ThreadPoolExecutor 的子类,重写 execute 方法,在提交任务的时候捕获当前线程的上下文快照,在任务真正执行之前恢复,执行完之后清理。阿里开源的 TransmittableThreadLocal(TTL)就是干这个的,它提供了 TtlExecutors 工具类,可以直接包装已有的线程池,改动成本很低。
// 用 TTL 包装线程池,上下文自动透传 Executor wrapped = TtlExecutors.getTtlExecutor(rpcPool); CompletableFuture.supplyAsync(() -> { // 这里能读到父线程 set 进去的 TraceId log.info("async task traceId={}", MDC.get("traceId")); return remoteCall(); }, wrapped);提示:包装线程池会带来额外的开销,主要体现在每次任务提交时的上下文拷贝。如果你的链路非常短、QPS 很高,要评估一下这部分开销。实测下来在常规业务量级下影响很小,通常在 1% 以内。
6.3 事务边界与异步混用的坑
数据库事务是基于 ThreadLocal 绑定连接实现的,异步线程拿不到主线程的事务上下文。这意味着在 @Transactional 方法里发起一个 supplyAsync,异步任务跑的是一套独立的事务,甚至可能是独立的数据源连接,写进去的数据主线程完全看不到。
更糟糕的情况是:主线程事务回滚了,但异步线程里执行的那部分数据库操作已经提交了。这时候数据就处于一种不一致的状态,而且非常难发现,因为两边都没有报错。
我的处理原则是:异步任务里不做写操作。所有的写库、更新状态、发消息这类动作,要么放在主线程的事务内同步执行,要么在事务提交之后再通过事件机制触发异步处理。如果确实需要在异步任务里写,那就让它用一个独立的事务,并且明确这是补偿逻辑,不是主流程的一部分。
读操作倒是可以在异步任务里放心做,因为它不涉及事务边界的问题。前面那个订单详情页的案例就是全读操作,所以适合并行。
6.4 任务堆积与线程池打满的排查路径
线上出现响应变慢、超时增多的时候,如果怀疑是线程池的问题,可以按这个顺序排查。
第一步,看监控面板上的活跃线程数和队列长度。活跃数打满且队列持续增长,基本可以确认是任务积压。
第二步,抓线程栈。用 jstack 或者 arthas 的 thread 命令,关注你那个线程池前缀的线程都在干什么。如果大量线程停在 socketRead、httpClient.execute、future.get 这类方法上,说明是下游变慢导致的阻塞。
第三步,确认是哪个下游。看线程栈里的调用链,或者结合上游的依赖监控,找出响应时间恶化的那个服务。这时候可以临时做两件事:调低该下游的超时时间,让它快速失败走降级;或者临时扩容线程池,先扛住流量。但扩容只是权宜之计,如果下游本身扛不住,扩线程池只会把压力传过去,问题从上游转移到下游。
第四步,复盘根因,做长期修复。要么是线程数配置不合理,要么是缺少超时和降级,要么是某个下游的容量规划出了问题。这三类问题对应的修复手段完全不同,别混在一起治。
还有一个容易被忽略的现象:如果拒绝策略用的是 CallerRunsPolicy,线程池满的时候任务会被提交线程自己执行,表现为"提交任务的接口变慢"。这时候你去看被提交的那个线程池,反而看不出什么异常,因为压力被转移到调用方了。排查的时候要注意这个陷阱,看看是不是 Tomcat 的工作线程在跑本应该在线程池里跑的任务。
7. 一些个人取舍经验
7.1 join 和 get 的选择
join 抛的是 CompletionException,get 抛的是 ExecutionException。前者是非受检异常,后者是受检异常。这个差别看起来很小,但在链式代码里影响很大。
如果链路上到处都是 get,你就得为每个调用写 try-catch,或者让方法签名抛出检查异常,代码会变得很难看。所以写编排的时候我一律用 join,只有在最外层需要把异常转换成业务异常的时候才用 get 或者统一 catch。
还有一个性能上的细微差别:在同一个 CompletableFuture 上多次调用 get 或者 join,它们的开销是可以忽略的,因为结果已经缓存了。所以不用担心在 allOf 的回调里连续 join 五个 future 会有性能问题。
7.2 不要在 thenApply 里做阻塞调用
这一条我踩过真实的坑。有一段代码是这样的:supplyAsync(fetch, pool).thenApply(this::enrich),enrich 里做了一次远程调用,没有加 Async,也没有加超时。结果就是,fetch 所在的线程执行完 fetch 之后接着执行 enrich,一阻塞就是几百毫秒。在 QPS 上来之后,线程池里的线程全被 enrich 占住,整个池子彻底瘫掉。
现在的规矩是:thenApply 里只放纯内存操作,任何涉及 I/O 的动作,要么用 thenApplyAsync 显式切换线程池,要么用 thenCompose 包一层新的 supplyAsync。判断标准就是问自己一句:这段代码会不会等待?会不会有不确定的耗时?只要答案是"可能会",就切出去。
7.3 保持编排的可读性比追求极致并行更重要
我见过为了把五个调用全塞进一条链里、结果写成三十行嵌套 lambda 的代码,也见过一个方法里拼接了七八种 API、后来的人根本不敢改的代码。异步编排的可读性和它的性能收益一样重要,因为代码是要被维护好几年的。
我的做法是控制单个链路的长度,超过五六个阶段就拆方法。把"取数据"和"组装数据"分开,取数据的方法返回一个包含若干个 CompletableFuture 的小容器对象,组装的方法只负责 join 和拼接。这样两部分的职责清晰,测试也好写。另外,给每个 CompletableFuture 变量起一个有业务含义的名字,比写成 future1、future2、future3 要重要得多,半年后回头看,名字就是最好的文档。
最后再分享一个小技巧。调试异步链路的时候,我习惯在每个关键阶段后面加一个 thenApply 打点,打印当前阶段名和时间戳。在预发环境打开,生产环境通过开关控制。这东西在排查"到底卡在哪一步"的时候非常有用,比看线程栈直观得多。而且它的成本很低,几行代码的事。
还有一个我一直在坚持的习惯:任何一条 CompletableFuture 链路的末端,必须有三个东西——异常处理、日志、耗时统计。少了任何一个,等线上出问题的时候你就得靠猜。这个习惯刚开始写的时候会觉得啰嗦,写顺了之后会发现它是你半夜能睡好觉的底气。