JAVA · Vol.II · DAY 14 · 并发编程
CompletableFuture 异步编排:从回调地狱到流式组合
开场:商品详情页为什么要 800ms?
某次大促压测,商品详情页接口 P99 高达 1.2 秒,是竞品的三倍。架构师打开链路监控一看:四路后端调用是串行跑的。为什么不能一起发出去?
一个商品详情页,页面渲染需要 4 个后端数据,它们之间完全没有依赖关系:
| 数据模块 | 调用方式 | 平均耗时 |
|---|---|---|
| 商品基础信息 | DB 查询 | 200ms |
| 价格与优惠 | 价格中心 RPC | 200ms |
| 用户评价 | 评价服务 RPC | 200ms |
| 推荐商品 | 推荐引擎 RPC | 200ms |
如果按顺序串行调用,总耗时 = 4 × 200ms = 800ms;而这 4 个调用彼此独立,完全可以在同一个时刻一起发出,等全部完成后再组装页面——这就是典型的「扇出 / 扇入(Fan-out / Fan-in)」。CompletableFuture 正是为这类「异步编排」而生的。
// 串行:800ms —— 前一个没回来,后一个不能出发 Product product = productService.getById(id); // 200ms Price price = priceService.getPrice(id); // 200ms List<Review> reviews = reviewService.list(id); // 200ms List<Recommend> recs = recommendService.list(id); // 200ms // 总计 ≈ 800ms // 并行 + 编排:≈ 300ms(四路同时发,等最慢的一路) CompletableFuture<Product> f1 = supplyAsync(() -> productService.getById(id)); CompletableFuture<Price> f2 = supplyAsync(() -> priceService.getPrice(id)); CompletableFuture<List<Review>> f3 = supplyAsync(() -> reviewService.list(id)); CompletableFuture<List<Recommend>> f4 = supplyAsync(() -> recommendService.list(id)); CompletableFuture.allOf(f1, f2, f3, f4).join(); // 等全部完成 ≈ 200ms DetailVO vo = assemble(f1.get(), f2.get(), f3.get(), f4.get()); // 组装 ≈ 100ms // 总计 ≈ 300ms —— RT 直接砍到原来的 1/3
CompletableFuture 是 JDK 8 引入的异步编排利器,面试官考它的原因:① 「多路 RPC 聚合」是生产里最常见的 IO 场景,做没做过、怎么做,一问便知;② 它串起了异步编程的全部知识点——线程池、回调(Callback)、异常传播、超时控制;③ 用得好能把接口 RT 砍到原来的 1/N,这是实打实的架构价值。
但在动手之前,先问一个问题:JDK 5 就有 Future 了,为什么还要等到 JDK 8 的 CompletableFuture?下一站,我们先看看传统 Future 到底差在哪。
传统 Future 的三大痛点:阻塞、无回调、不能组合
JDK 5 引入的 Future 已经能做到「任务丢到线程池、之后取结果」,但用它写「多个异步任务的编排」非常别扭。先看代码,再总结痛点:
ExecutorService pool = Executors.newFixedThreadPool(4); Future<Product> f1 = pool.submit(() -> productService.getById(id)); Future<Price> f2 = pool.submit(() -> priceService.getPrice(id)); // ① get() 会阻塞调用线程 —— 无论 f2 多快,只要 f1 慢,这条线程就干等着 Product p = f1.get(); // 阻塞,最多等 800ms,期间线程啥也干不了 // ② 想"结果一出来就处理"?没有回调,只能轮询 isDone() while (!f2.isDone()) { // 轮询浪费 CPU,且延迟不可控 Thread.sleep(10); } // ③ 想组合?f3 依赖 f1 和 f2 的结果 —— 只能先 get 出来,再手动 submit 一个新任务 Price pr = f2.get(); // 取结果的顺序被代码写死,无法"谁先好先处理谁" Future<Coupon> f3 = pool.submit(() -> couponService.get(p, pr));
三个痛点点名:
| 痛点 | Future 的表现 | 带来的问题 |
|---|---|---|
| 阻塞 | get() 会挂起调用线程 | Tomcat 线程池里一条线程被白白占住;本该处理 20 个请求的时间只处理了 1 个 |
| 无回调 | 结果就绪时不会通知你 | 只能阻塞 get() 或轮询 isDone(),前者占线程、后者耗 CPU 且不及时 |
| 不能组合 | 没有「合并两个 Future」「等全部完成」的 API | 依赖关系要手写:先 get 再 submit,代码变成一串脆弱的串行步骤;等全部完成只能自己上 CountDownLatch |
阻塞 = 站在店门口干等,啥也别干;轮询 = 每 5 分钟刷一次手机,又累又慢还可能错过;回调 = 下单时留下手机号,外卖到了骑手自动打电话通知你——你的时间完全解放。Future 只支持前两种,CompletableFuture 才支持第三种。
面试追问:那第 1 站用 4 个 Future + 4 次 get() 不也能并行吗?为什么说它不行?
能并行、不能编排。4 个 Future 一起 submit,确实能同时跑;但取结果的代码必须按固定顺序 f1.get() → f2.get() → …,如果 f1 是最慢的那路,即使 f4 早就好了,你也得等 f1 结束才能走到下一行——「谁先好先处理谁」做不到,更别说「f2 的结果要喂给 f3」这种依赖链。另外 Future.get() 抛的是受检异常(InterruptedException + ExecutionException),每个调用点都要 try-catch,代码里全是噪音。
回调地狱:嵌套回调为什么让人抓狂
既然 Future 没有回调,当年没有 CompletableFuture 的开发者是怎么解决「结果就绪自动处理」的?答案是回调(Callback)——把「下一步做什么」作为参数传给异步方法。但回调一旦嵌套,就成了著名的回调地狱(Callback Hell):
// 查用户 → 查会员价 → 查优惠券,每步都依赖上一步 ListenableFuture<User> userF = executor.submit(() -> userService.get(userId)); Futures.addCallback(userF, new FutureCallback<User>() { public void onSuccess(User user) { // 第 1 层 ListenableFuture<Price> priceF = executor.submit(() -> priceService.get(user)); Futures.addCallback(priceF, new FutureCallback<Price>() { public void onSuccess(Price price) { // 第 2 层 ListenableFuture<Coupon> couponF = executor.submit(() -> couponService.get(price)); Futures.addCallback(couponF, new FutureCallback<Coupon>() { public void onSuccess(Coupon coupon) { // 第 3 层,才拿到全部数据 render(user, price, coupon); } public void onFailure(Throwable t) { // 每层都要写一遍错误处理 log.error("查询优惠券失败", t); } }, executor); } public void onFailure(Throwable t) { log.error("查询价格失败", t); } }, executor); } public void onFailure(Throwable t) { log.error("查询用户失败", t); } }, executor);
这段代码的问题肉眼可见:
- 横向膨胀:每多一个依赖步骤,代码就往右缩进一层,三层之后可读性归零——业务逻辑被「括号」淹没;
- 错误处理重复:每层都要单独写一遍
onFailure,漏一层异常就静默丢失; - 依赖关系不可见:谁依赖谁被埋在嵌套结构里,想抽出来复用、想调整顺序,牵一发动全身。
每打开一层,发现里面还有一层。回调地狱的本质是:把「流程」写成了「嵌套的盒子」。CompletableFuture 的答案是把盒子拆开、拍平成一维的链——这个对比我们留到第 14 站,同一段需求两种写法摆在一起看,你会立刻明白它解决了什么。
CompletableFuture 是什么:Future + 完成即触发
CompletableFuture 同时实现了两个接口:Future(能拿结果、能取消)+ CompletionStage(完成阶段:能在「完成」这个事件上挂一串后续动作)。它的核心机制一句话——完成即触发(Completion):任务完成的那一刻,所有提前登记好的依赖动作会被自动逐个触发,不需要你手动通知、不需要轮询。
内部实现上(JDK 8),它就是两个字段:一个 volatile Object result 存结果(或异常),一个 volatile Completion stack 存所有依赖动作——用无锁的 Treiber 栈组织,谁调 complete(),谁负责把这个栈弹出来逐个执行。所以「注册回调」和「完成任务」完全解耦,可以发生在任何线程、任何顺序:
CompletableFuture<String> f = new CompletableFuture<>(); // 手动创建,此刻还没有任务 f.thenAccept(result -> log.info("收到结果:{}", result)); // ① 先"登记"依赖动作 new Thread(() -> { // ② 另一个线程 1 秒后完成 Thread.sleep(1000); f.complete("hello"); // ③ 触发上面的回调(幂等,只有第一次生效) }).start(); log.info("主线程继续做自己的事"); // ④ 主线程不用等、不用轮询 // 输出顺序:先"主线程继续做自己的事",1 秒后"收到结果:hello"
这也是 complete() / completeExceptionally() 的意义:CompletableFuture 不一定要绑一个任务,它可以被任何线程在任何时刻手动完成——第 12 站的超时降级,本质就是「定时器线程帮它 complete」。三个接口的关系一张表看清:
| 维度 | Future(JDK 5) | CompletableFuture(JDK 8) | CompletionStage(JDK 8) |
|---|---|---|---|
| 角色 | 异步结果的占位符 | Future + CompletionStage 的实现类 | 「完成阶段」接口:定义各种编排方法 |
| 回调 | ❌ 无 | ✅ thenXxx 系列 | ✅(接口方法) |
| 组合 | ❌ 无 | ✅ thenCombine / allOf / anyOf | ✅ |
| 手动完成 | ❌ 不能 | ✅ complete() / completeExceptionally() | — |
| 线程池控制 | 提交时指定 | 创建 + 每个编排动作都能指定 | — |
创建 CompletableFuture:三种姿势 + 手动完成
创建 CompletableFuture 有三种常用方式,核心区别在于是否有返回值以及是否已经完成:
// ① supplyAsync —— 有返回值(Supplier → CompletableFuture<T>) CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> { // 耗时操作,最终返回结果 return productService.getById(123).getName(); }); // ② runAsync —— 无返回值(Runnable → CompletableFuture<Void>) CompletableFuture<Void> f2 = CompletableFuture.runAsync(() -> { // 只执行操作,不返回结果(如发通知、写日志、埋点) notifyService.send("order_created"); }); // ③ completedFuture —— 已经完成的 Future(测试、默认值、短路) CompletableFuture<String> f3 = CompletableFuture.completedFuture("cached-value"); f3.isDone(); // true,立即完成,后面接 thenXxx 会立刻执行
除此之外,还有两个「手动完成」的方法——上一站见过 complete(),与之对应的是 completeExceptionally(Throwable):把异常作为完成结果。注意 complete() 的返回值:true 表示这次调用真的完成了它(第一次),false 表示它早已完成、本次调用被忽略——幂等。它有个很实用的应用:本地缓存命中时,用 completedFuture 短路,下游完全不用感知「这次是缓存还是真的查了」:
CompletableFuture<Price> loadPrice(long id) { Price cached = localCache.get(id); if (cached != null) { return CompletableFuture.completedFuture(cached); // 命中缓存:立即完成 } return CompletableFuture.supplyAsync(() -> priceService.getPrice(id), pool); }
| 方式 | 入参 | 返回 | 是否异步执行 | 典型场景 |
|---|---|---|---|---|
supplyAsync | Supplier<T> | CompletableFuture<T> | 是 | 有返回值的异步任务(查 DB、调 RPC) |
runAsync | Runnable | CompletableFuture<Void> | 是 | 无返回值的异步动作(发通知、写日志) |
completedFuture | T | CompletableFuture<T> | 否(立即完成) | 测试替身、缓存短路、默认值 |
new + complete() | — | — | 否 | 手动控制完成时机(配合定时器做超时) |
supplyAsync 有返回值,runAsync 无返回值,completedFuture 立即可用。补充:两个 Async 创建方法都有带 Executor 的重载——生产环境必须用这个重载,原因下一站讲。
线程池选择:为什么生产必须传自定义线程池
supplyAsync / runAsync 不传 Executor 时,默认使用 ForkJoinPool.commonPool()——这是全 JVM 共享的一个公共线程池,线程数 = CPU 核数 - 1(可用 -Djava.util.concurrent.ForkJoinPool.common.parallelism 覆盖)。听起来很方便,生产里却是个大坑:
// 4 核机器 → commonPool 只有 3 个线程 CompletableFuture.supplyAsync(() -> rpcA.call()); // 默认丢进 commonPool,等 RPC 200ms CompletableFuture.supplyAsync(() -> rpcB.call()); // 又占 1 个 CompletableFuture.supplyAsync(() -> rpcC.call()); // 又占 1 个 —— 3 个全在阻塞等 IO // 此时全项目的 parallelStream() 和所有不传线程池的 CF 一起饿死排队! list.parallelStream().map(...).collect(...); // ← 也被拖垮
RPC / DB 这类 IO 密集任务的特点是「线程大部分时间在等」,等的时候线程池线程被占着不放。commonPool 只有核数 - 1 条线程,一旦被几个慢 RPC 占满,你的 parallelStream、别人模块的异步任务、甚至 JDK 内部依赖 commonPool 的机制全部一起遭殃——一个服务的异步代码,把全 JVM 的公共池污染了。这就像公共泳池:一个人抽筋,全池子的人都得绕着走。
ExecutorService bizPool = new ThreadPoolExecutor(16, 32, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(1000), new ThreadFactoryBuilder().setNameFormat("biz-async-%d").build(), new ThreadPoolExecutor.CallerRunsPolicy()); CompletableFuture<Product> f = CompletableFuture.supplyAsync( () -> productService.getById(123), bizPool); // ← 第二个参数指定线程池
两个生产要点:
- 线程池按业务隔离:查询聚合、异步写、定时任务各用各的池子,参数按各自场景调(IO 密集适当放大,参考卷二线程池篇的
N*(1+W/C)公式),互相不拖累、监控也清晰; - 拒绝策略的坑:示例里用了
CallerRunsPolicy——队列满了会在「提交任务的那条线程」上直接跑。如果提交方是请求线程,一个慢任务可能把请求线程也拖住,所以必须配合超时兜底(第 12 站);对异步链路而言,也可以考虑丢弃 + 降级。
追问:thenApply 这类回调到底跑在哪条线程上?
不带 Async 后缀的方法(thenApply / thenAccept 等)在「触发完成的那条线程」上同步执行——谁 complete 的,谁顺手把回调也跑了;带 Async 后缀的方法(thenApplyAsync 等)才丢进线程池。所以:回调里如果要做耗时操作,要么用 Async 变体、要么保证完成线程不阻塞(比如别在 RPC 的 IO 线程里做重活)。这也是很多线上 bug 的来源:回调在 IO 线程里跑了 500ms 的 CPU 计算,把下游连接池拖垮。
生产环境必须传自定义 Executor,不要用默认的 ForkJoinPool.commonPool()——线程数不可控、全 JVM 共享、互相干扰。Async 变体换线程池执行,非 Async 变体在完成线程上执行。
链式转换:thenApply / thenAccept / thenRun
CompletableFuture 最强大的特性是链式组合——像 Stream 一样把多个异步步骤串成流水线(Pipeline)。面试最爱问的三个方法:thenApply(转换)、thenAccept(消费)、thenRun(执行动作)。
// thenApply:同步转换(Function<T, U>),类比 Stream.map():输入 T,输出 U CompletableFuture<String> nameFuture = supplyAsync(() -> productService.getById(123)) // CompletableFuture<Product> .thenApply(product -> product.getName()); // CompletableFuture<String> // thenApplyAsync:转换本身是耗时操作时,丢进线程池异步执行 f.thenApply(result -> result.toUpperCase()); // 在完成线程上同步执行(快操作) f.thenApplyAsync(result -> callRPC(result)); // 提交到线程池执行(慢操作) // thenAccept:消费结果,无返回值(类比 forEach) supplyAsync(() -> orderService.create(dto)) .thenAccept(orderId -> log.info("订单创建成功: {}", orderId)); // thenRun:不关心上一步结果,只执行后续动作(上一步结果被丢弃) supplyAsync(() -> orderService.create(dto)) .thenRun(() -> metricsService.increment("order.created"));
thenApply ≈ map(T → U,有返回值) | thenAccept ≈ forEach(T → void,消费) | thenRun(跟结果无关,跑个动作)
选型口诀:要转换结果用 thenApply,要消费结果用 thenAccept,跟结果没关系用 thenRun。方法名后缀规律也要记住:Async 表示丢线程池、不带则同步执行(执行线程的语义见第 6 站)。
thenApply 有参有返回、thenAccept 有参无返回、thenRun 无参无返回。它们都只适合「下一步是同步操作」——如果下一步本身是异步的(返回 CompletableFuture),就要用下一站的 thenCompose。
thenCompose vs thenApply:嵌套与拍平
这是面试出现频率最高的追问:「thenApply 和 thenCompose 有什么区别?」一句话版:thenApply 的回调返回普通值 U,thenCompose 的回调返回 CompletableFuture<U>——thenCompose 会把返回的 Future 自动「拍平」,避免嵌套。
// thenApply:同步转换(Function<T, U>)—— 回调里做的是同步操作 CompletableFuture<String> nameF = supplyAsync(() -> productService.getById(123)) // CompletableFuture<Product> .thenApply(product -> product.getName()); // getName() 是同步方法 → CompletableFuture<String> // thenCompose:扁平化转换(Function<T, CompletableFuture<U>>)—— 回调里做的是异步操作 CompletableFuture<Price> priceF = supplyAsync(() -> productService.getById(123)) // CompletableFuture<Product> .thenCompose(product -> priceService.getPriceAsync(product.getSkuId())); // 返回 CompletableFuture<Price>,自动拍平
如果用 thenApply 去接一个异步操作,会发生什么?嵌套——结果变成 CompletableFuture<CompletableFuture<Price>>,取结果要多等一层:
// ✗ 错误:用 thenApply 串联异步步骤,得到套娃 CompletableFuture<CompletableFuture<Price>> nested = supplyAsync(() -> productService.getById(123)) .thenApply(product -> priceService.getPriceAsync(product.getSkuId())); // 嵌套! Price price = nested.join().join(); // 只能再 join 一层 —— 又阻塞又难看 // ✓ 正确:thenCompose 自动拍平,链上始终只有一层 CompletableFuture<Price> flat = supplyAsync(() -> productService.getById(123)) .thenCompose(product -> priceService.getPriceAsync(product.getSkuId())); Price price = flat.join(); // 一层搞定
thenCompose 就是异步版的 flatMap:flatMap 把「元素 → 集合」拍平成一层流,thenCompose 把「值 → Future」拍平成一层 Future。判断口诀:回调里是同步方法用 thenApply,回调里返回 Future 用 thenCompose——永远不要让链上出现 CompletableFuture<CompletableFuture<…>>。
thenApply 回调返回普通值(map),thenCompose 回调返回 Future 并自动拍平(flatMap)。下一步是异步操作用 thenCompose,否则会得到嵌套的 CompletableFuture,取结果被迫二次阻塞。
组合与聚合:thenCombine / allOf / anyOf
链式转换解决的是「串行流水线」,但真实场景中经常需要多个独立 Future 的聚合。这就是 thenCombine、allOf、anyOf 的舞台。
CompletableFuture<Product> productF = supplyAsync(() -> productService.getById(id)); CompletableFuture<Price> priceF = supplyAsync(() -> priceService.getPrice(id)); CompletableFuture<DetailVO> detailF = productF.thenCombine(priceF, (product, price) -> { DetailVO vo = new DetailVO(); vo.setName(product.getName()); vo.setPrice(price.getAmount()); return vo; });
CompletableFuture<Void> all = CompletableFuture.allOf(f1, f2, f3, f4); all.thenRun(() -> { // 所有 Future 已完成,安全 join() DetailVO vo = assemble(f1.join(), f2.join(), f3.join(), f4.join()); response.complete(vo); }); // 注意:allOf 返回 CompletableFuture<Void>,不携带结果,必须手动 join 每个子 Future
// 多数据源竞速:谁先返回用谁(常用于缓存 + DB 双读) CompletableFuture<String> fromCache = supplyAsync(() -> cacheService.get(key)); CompletableFuture<String> fromDB = supplyAsync(() -> dbService.query(key)); CompletableFuture<Object> fastest = CompletableFuture.anyOf(fromCache, fromDB); // 注意:anyOf 返回 CompletableFuture<Object>,需要强转
① allOf 的失败语义:只要其中一个 Future 异常完成,返回的 all 也异常完成——所以 all.thenRun(...) 里的回调不会执行;但其它已经成功完成的任务不受影响,各自的 join() 仍然能取到结果(谁失败谁抛)。② anyOf 的应用:不止「缓存 vs DB」双读,多机房容灾「取最快到达的那路」也是它——但注意要配合超时,否则最慢那路永远不完成,anyOf 本身不设超时也会一直等。
thenCombine 是「两路合并」,allOf 是「等所有」,anyOf 是「抢最快」。allOf 返回 Void 需要手动取结果,anyOf 返回 Object 需要强转。
异常处理三兄弟:exceptionally / handle / whenComplete
异步编程最容易被忽视的就是异常处理。CompletableFuture 中的异常不会抛到调用方线程——如果你不主动处理,它会静默吞掉:异常被存进 result 字段(异常完成),后续依赖动作全部跳过,直到你调 get() / join() 时才以 CompletionException 的形式爆发。所以每条链都要有「兜底」。
// ① exceptionally —— 异常兜底,返回同类型默认值(类比 catch + return 默认值) CompletableFuture<List<Review>> reviewsF = supplyAsync(() -> reviewService.list(productId)) .exceptionally(ex -> { log.warn("评价服务异常", ex); return Collections.emptyList(); // ← 必须返回同类型 }); // ② handle —— 成功和异常都处理,可返回不同类型(类比 try-catch-finally 里统一加工) CompletableFuture<String> resultF = supplyAsync(() -> riskyCall()) .handle((result, ex) -> ex != null ? "fallback" : result); // ③ whenComplete —— 只读副作用(日志/埋点),不改变结果(类比 finally) supplyAsync(() -> orderService.create(dto)) .whenComplete((orderId, ex) -> { if (ex != null) log.error("创建订单失败", ex); else log.info("创建订单成功: {}", orderId); }); // 结果原样往下传,异常也原样往下传
| 方法 | 输入 | 输出 | 是否改变结果 | 典型场景 |
|---|---|---|---|---|
exceptionally | Throwable | 同类型 T | 是(兜底值) | 降级返回默认值 |
handle | (T, Throwable) | 新类型 U | 是 | 成功/失败都需要转换 |
whenComplete | (T, Throwable) | 同类型 T | 否 | 日志、埋点、监控 |
exceptionally 的回调必须返回与原 Future 相同泛型类型的值。新手常犯的错误是在 exceptionally 中写 return null,导致下游收到 null 引发 NPE。正确做法是返回有意义的默认值(如空列表、空对象),或者干脆在源头用 completeOnTimeout 兜底(第 12 站)。
// 异常会跳过中间步骤,直接传到最近的 exceptionally supplyAsync(() -> { throw new RuntimeException("boom"); }) .thenApply(s -> s.toUpperCase()) // ← 跳过,不执行 .thenCompose(s -> callRPC(s)) // ← 跳过,不执行 .exceptionally(ex -> { log.error("链路异常", ex); return "default"; });
这张图要刻在脑子里:异常沿依赖链向下传播,中间步骤全部短路,直到遇到 exceptionally / handle 才被接住;如果一路都没有处理者,异常就「沉底」——看起来任务静静结束了,其实是个哑弹,等 get/join 才炸。
exceptionally 是兜底降级(必须返回同类型值),handle 是成功/失败双通道处理(可换类型),whenComplete 是只读副作用。链式调用中异常会「穿透」中间步骤,直到被 exceptionally 或 handle 捕获;没人接住就会沉底,等到 get/join 才爆发。
get() vs join():受检异常与运行时异常
get() 和 join() 都能「阻塞地取结果」,区别全在异常处理方式上,这也是面试爱问的细节:
// get():抛两个受检异常,编译器逼你处理 try { Product p = f1.get(); // throws InterruptedException, ExecutionException } catch (InterruptedException e) { // 等待期间线程被中断 Thread.currentThread().interrupt(); // 最佳实践:恢复中断标志 } catch (ExecutionException e) { // 任务执行失败 Throwable cause = e.getCause(); // 原始异常在这里(要剥一层) log.error("任务执行失败", cause); } // join():只抛 CompletionException(运行时异常),不写 try 也能编译 Product p = f2.join(); // 失败了直接抛出来,自己接住 // cause 同样在 e.getCause() 里,异常类型是 CompletionException
| 维度 | get() | join() |
|---|---|---|
| 抛出的异常 | InterruptedException + ExecutionException | CompletionException |
| 是否受检 | 受检(checked),必须 try-catch | 非受检(unchecked),可不处理 |
| 中断语义 | 可中断:等待中被 interrupt 会抛异常 | 不可中断:中断只被记下,等完再恢复标志 |
| 异常包装 | 原始异常在 ExecutionException.getCause() | 原始异常在 CompletionException.getCause() |
| 典型场景 | 最外层 Controller 取最终结果,显式处理异常 | 链式内部、测试代码,追求简洁 |
还有一个不阻塞的变体值得记住——getNow(valueIfAbsent):Future 已完成就返回结果,未完成就直接返回你给的默认值,绝不阻塞。适合「尽力而为」的读取:
// 缓存异步预热中:读到了就用,没读到先返回默认值,绝不阻塞请求线程 CompletableFuture<Price> preheating = loadPrice(id); // 可能还没完成 Price price = preheating.getNow(Price.UNKNOWN); // 未完成 → 立即返回 UNKNOWN
get() 抛受检异常(InterruptedException + ExecutionException),join() 抛非受检的 CompletionException;两者都把原始异常包在 getCause() 里。链路内部建议统一 join(),最外层用 get() 显式兜底;需要「不阻塞拿结果」用 getNow()。
超时控制:orTimeout / completeOnTimeout 与 JDK 8 方案
生产环境中,异步调用必须有超时——否则一个下游服务 hang 住,你的线程会被无限阻塞,最终线程池耗尽、服务雪崩。没有超时的异步调用,就像挂机等一个永远不会来的电话。
// orTimeout —— 超时则把 future 标记为"异常完成"(TimeoutException) CompletableFuture<Product> f = supplyAsync(() -> productService.getById(id)) .orTimeout(500, TimeUnit.MILLISECONDS); // 500ms 内未完成 → future 以 TimeoutException 完成 // 后续 exceptionally / handle 可以捕获该异常做降级 // completeOnTimeout —— 超时则用默认值"正常完成"(不会抛异常) CompletableFuture<List<Review>> reviewsF = supplyAsync(() -> reviewService.list(id)) .completeOnTimeout(Collections.emptyList(), 300, TimeUnit.MILLISECONDS); // 300ms 内未完成 → 自动用空列表完成,下游正常执行
| 方法 | 超时后行为 | 适用场景 |
|---|---|---|
orTimeout | 以 TimeoutException 异常完成 | 超时需要报错 / 触发 exceptionally 做复杂兜底 |
completeOnTimeout | 以指定默认值正常完成 | 超时直接降级返回默认值(评价、推荐等非核心模块) |
一个小细节:orTimeout / completeOnTimeout 内部靠一个共享的静态 ScheduledThreadPoolExecutor(Delayer)定时「叫醒」——本质还是第 4 站说的「定时器线程调 complete() / completeExceptionally()」,幂等性保证了它不会覆盖真正的结果。
// JDK 8 没有 orTimeout / completeOnTimeout,需手动实现(原理一模一样) private static final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(2, r -> { Thread t = new Thread(r, "cf-timeout-scheduler"); t.setDaemon(true); return t; }); public static <T> CompletableFuture<T> withTimeout( CompletableFuture<T> future, long timeout, TimeUnit unit) { scheduler.schedule(() -> future.completeExceptionally(new TimeoutException("timed out")), timeout, unit); return future; // complete 幂等,Future 已完成则此调用无效 } CompletableFuture<Product> f = withTimeout( supplyAsync(() -> productService.getById(id)), 500, TimeUnit.MILLISECONDS);
关键点:complete() 和 completeExceptionally() 都是幂等的——只有第一次调用生效。如果 Future 已经正常完成,后续的超时回调调 completeExceptionally 会被忽略,这让手动超时方案非常安全。
JDK 9+ 用 orTimeout(超时异常)和 completeOnTimeout(超时降级)。JDK 8 用 ScheduledExecutorService + completeExceptionally 手动实现。生产环境每个异步调用都必须设超时——这是防止线程池耗尽雪崩的底线。
生产实战:商品详情页完整编排
面试中如果能写出以下完整模式,会非常有加分。我们用商品详情页把三种核心编排模式串起来(以下模板假设 JDK 9+;JDK 8 请把 orTimeout / completeOnTimeout 换成第 12 站的 withTimeout 工具方法):
public CompletableFuture<DetailVO> getProductDetail(long productId) { // Pipeline:查用户 → 查会员价(有依赖,用 thenCompose 串联) CompletableFuture<PriceVO> priceF = supplyAsync(() -> userService.getCurrentUser(), pool) .thenCompose(user -> priceService.getMemberPriceAsync(productId, user.getLevel())) .orTimeout(300, TimeUnit.MILLISECONDS) .exceptionally(ex -> priceService.getDefaultPrice(productId)); // Fan-out:商品、评价、推荐无依赖,并行发起 CompletableFuture<Product> productF = supplyAsync(() -> productService.getById(productId), pool) .orTimeout(200, TimeUnit.MILLISECONDS); CompletableFuture<List<Review>> reviewsF = supplyAsync(() -> reviewService.list(productId), pool) .completeOnTimeout(Collections.emptyList(), 300, TimeUnit.MILLISECONDS); CompletableFuture<List<Recommend>> recsF = supplyAsync(() -> recommendService.list(productId), pool) .completeOnTimeout(Collections.emptyList(), 300, TimeUnit.MILLISECONDS); // Fan-in:聚合所有结果 return CompletableFuture.allOf(productF, priceF, reviewsF, recsF) .thenApply(v -> { DetailVO vo = new DetailVO(); vo.setProduct(productF.join()); vo.setPrice(priceF.join()); vo.setReviews(reviewsF.join()); vo.setRecommendations(recsF.join()); return vo; }); }
public static <T> CompletableFuture<T> withRetry( Supplier<CompletableFuture<T>> action, int maxRetries, long timeoutMs) { return action.get() .orTimeout(timeoutMs, TimeUnit.MILLISECONDS) .exceptionally(ex -> { if (maxRetries > 0) return withRetry(action, maxRetries - 1, timeoutMs).join(); throw new CompletionException("重试耗尽", ex); }); } // 使用:最多重试 2 次,每次超时 500ms CompletableFuture<Product> f = withRetry( () -> supplyAsync(() -> productService.getById(id), pool), 2, 500);
| 模式 | 特点 | 适用场景 |
|---|---|---|
| Pipeline | thenCompose 串联,前一步输出是后一步输入 | 有依赖的多步异步(如:查用户 → 查会员价) |
| Fan-out / Fan-in | allOf 并行发起,等全部完成再聚合 | 无依赖的多路聚合(如:商品 + 评价 + 推荐) |
| Timeout + Retry | orTimeout + exceptionally 递归重试 | 不稳定的下游调用,需要容错 |
推荐是非核心模块。用 completeOnTimeout 超时降级为空列表 + exceptionally 捕获异常也返回空列表。核心模块(商品、价格)才需要重试。原则:核心模块重试 + 超时,非核心模块降级 + 超时——重试也要设上限,且最好退避,否则雪崩时你的重试就是在补刀。
CF 的回调经常在别的线程上执行(线程池线程、完成线程),而日志框架的 traceId 通常放在 ThreadLocal 里——跨线程就丢了,全链路追踪在异步段断链,排查问题抓瞎。解决:① 提交任务前把 traceId 抓出来,回调里手动塞回;② 用阿里开源的 TransmittableThreadLocal 做线程池透传(卷二 ThreadLocal 篇有专题)。
回调地狱 vs 流式编排:同一个需求两种写法
把第 3 站那个「查用户 → 查会员价 → 算实付」的需求,用两种方式各写一遍。先看对比,再谈感受:
// 每一步都要:传回调 + 传错误处理,越嵌越深
fetchUser(userId, user ->
fetchMemberPrice(user, price ->
calcPay(price, pay -> render(pay), onError),
onError),
onError);
CompletableFuture.supplyAsync(() -> userService.get(userId), pool)
.thenCompose(user -> priceService.getMemberPriceAsync(user, pool))
.thenApply(price -> calcService.calcPay(price))
.exceptionally(ex -> defaultPay())
.thenAccept(pay -> render(pay));
| 维度 | 回调嵌套 | 流式编排 |
|---|---|---|
| 代码形态 | 横向膨胀,括号层层包裹 | 一维链,每步一行,从上往下读就是流程 |
| 异常处理 | 每层单独写一遍 onFailure | 链尾一个 exceptionally 统一兜底 |
| 组合能力 | 手动拼,改一步动全链 | thenCombine / allOf / anyOf 随意拼装 |
| 线程控制 | 靠外部 executor 参数传递 | 每步可指定 Async + 线程池 |
| 调试排查 | 回调栈信息支离破碎 | 链式结构清晰,异常 cause 可追溯 |
回调地狱的问题是「流程」被写成了「嵌套」;CompletableFuture 把嵌套拍平,让数据沿着一维的链流动——上一站的输出自动成为下一站的输入,异常自动向下传播。这就是「从回调地狱到流式编排(Reactive-style Chaining)」的全部意义:不是少了几个括号,而是让流程可读、可复用、可组合。
方法速查表:一张表记住全部 API
面试前把这张表过一遍,主要 API 就不会漏:
// 还没列全的"双 Future 变体":thenAcceptBoth(两个都完成再消费)、 // applyToEither / acceptEither(两个任一完成就处理)。面试提到"还有变体"即可加分。
| 类别 | 方法 | 一句话说明 | Stream 类比 |
|---|---|---|---|
| 创建 | supplyAsync | 有返回值的异步任务 | — |
runAsync | 无返回值的异步任务 | — | |
completedFuture | 已完成的 Future(测试/默认值/短路) | Stream.of() | |
| 链式转换 | thenApply | 同步转换 T → U | map |
thenCompose | 异步扁平转换 T → CF<U>,自动拍平 | flatMap | |
thenAccept / thenRun | 消费结果 / 执行动作 | forEach | |
| 组合聚合 | thenCombine | 两个 Future 合并为一个结果 | — |
allOf | 等所有 Future 完成(返回 Void) | — | |
anyOf | 最先完成的胜出(返回 Object) | — | |
| 异常处理 | exceptionally | 异常兜底,返回同类型默认值 | — |
handle | 成功/失败双通道处理,可换类型 | — | |
whenComplete | 只读副作用(日志/埋点),结果不变 | peek | |
| 结果获取 | get() / join() | 阻塞取结果(受检 / 非受检异常) | — |
getNow(default) | 已完成返回结果,未完成返回默认值,不阻塞 | — | |
| 手动完成 | complete(v) | 用值完成(幂等) | — |
completeExceptionally(t) | 用异常完成(幂等) | — | |
| 超时 (9+) | orTimeout | 超时 → TimeoutException 异常完成 | — |
completeOnTimeout | 超时 → 用默认值正常完成 | — |
后缀决定一切:Apply 有参有返回、Accept 有参无返回、Run 无参无返回、Compose 拍平嵌套;带 Async 丢线程池;Exceptionally / Handle 处理异常、WhenComplete 只读副作用;Combine / AllOf / AnyOf 处理多 Future。
常见坑清单与最佳实践
- 忘记传自定义线程池——默认 ForkJoinPool.commonPool() 线程数 = CPU 核数 - 1,IO 密集场景必炸,还拖垮别人的 parallelStream
- allOf 返回 Void——不能直接从 allOf 拿结果,必须手动 join() 每个子 Future
- exceptionally 返回 null——下游收到 null 引发 NPE,应返回有意义的默认值
- 没有设超时——下游 hang 住导致线程永久阻塞,最终线程池耗尽、服务雪崩
- thenApply vs thenCompose 混用——下一步是异步操作用 thenCompose,否则得到嵌套 Future,取结果二次阻塞
- join() / get() 选错——get() 是受检异常必须处理;链路内部统一 join() 更干净
- 在非 Async 回调里做耗时操作——回调在完成线程上跑,会阻塞那条线程(可能是 RPC 的 IO 线程)
- ThreadLocal 不传递——回调跨线程,traceId 丢失、全链路追踪断链(用 TransmittableThreadLocal)
- 异常被吞——没接住也没 get/join 的异常静默沉底,线上只剩"结果不对"的谜案;记得打日志或上报
- 重试无上限 / 无退避——雪崩时疯狂重试等于补刀,重试必须设上限 + 退避 + 熔断
// 团队内统一封装,禁止裸写 supplyAsync public class Asyncs { public static <T> CompletableFuture<T> supply(Supplier<T> task, long timeoutMs, Supplier<T> fallback) { return CompletableFuture.supplyAsync(task, BizPools.ASYNC) .completeOnTimeout(fallback.get(), timeoutMs, TimeUnit.MILLISECONDS) .exceptionally(ex -> { log.warn("async task failed, fallback", ex); return fallback.get(); }); } } // 调用方一行搞定:Asyncs.supply(() -> rpc.call(), 500, () -> defaultV)
监控也别忘:聚合接口要对每路子调用打点(耗时、成功率、超时率),谁慢了、谁抖了,一眼可见——异步化之后,单路延迟被并行掩盖,反而更需要逐路监控,否则慢接口会悄悄藏进 P99 里。
- 每个
supplyAsync都传自定义线程池,线程池按业务隔离、命名规范; - 每个异步调用都设超时:核心模块 orTimeout + 重试(有上限、有退避),非核心模块 completeOnTimeout 降级;
- 链路末端必须有异常处理者(exceptionally / handle),至少打日志;
- 回调里不做耗时操作,需要就上 Async 变体;
- 跨线程传值用 TransmittableThreadLocal,保住 traceId。
这一篇你掌握了什么
核心知识点回顾
- 为什么需要它:传统 Future 三大痛点——get() 阻塞、无回调(只能轮询)、不能组合;回调嵌套则形成回调地狱(Callback Hell),流程被埋进括号里;
- 核心模型:CompletableFuture = Future + CompletionStage,完成即触发(Completion)——结果是字段、依赖动作是栈,谁 complete 谁唤醒;complete() / completeExceptionally() 幂等,可被任何线程在任何时刻调用;
- 创建:supplyAsync(有返回值)/ runAsync(无返回值)/ completedFuture(立即可用,缓存短路、测试);
- 线程池:默认 ForkJoinPool.commonPool()(核数 - 1、全 JVM 共享),生产必须传自定义线程池;非 Async 回调在完成线程执行、Async 变体才换线程池;
- 链式转换:thenApply(map,有参有返回)/ thenAccept(有参无返回)/ thenRun(无参);下一步是异步操作用 thenCompose(flatMap)拍平,避免 CompletableFuture<CompletableFuture<…>>;
- 组合聚合:thenCombine(两路合并)、allOf(等全部,返回 Void 需手动 join)、anyOf(抢最快,返回 Object 需强转);
- 异常处理:exceptionally(兜底同类型值)/ handle(成败双通道可换类型)/ whenComplete(只读副作用);异常沿依赖链穿透,没人接住就沉底,get/join 时才爆发;
- 取结果:get() 抛受检 InterruptedException + ExecutionException,join() 抛非受检 CompletionException;getNow() 不阻塞;
- 超时控制:JDK 9+ 用 orTimeout(异常完成)/ completeOnTimeout(默认值完成),JDK 8 用 ScheduledExecutorService 手动实现;每个异步调用都必须设超时;
- 生产模式:Pipeline(thenCompose 串依赖)+ Fan-out/Fan-in(allOf 并行聚合)+ Timeout+Retry(核心模块重试、非核心模块降级);注意 traceId 跨线程透传。
一句话总结
CompletableFuture 把「阻塞等待」变成「完成即触发」,把「回调嵌套」拍平成「一维流式链」:创建(supplyAsync)→ 转换(thenApply/thenCompose)→ 聚合(thenCombine/allOf/anyOf)→ 兜底(exceptionally/handle)→ 超时(orTimeout/completeOnTimeout),配合自定义线程池,就能写出 RT 更低、可读性更高、可维护的异步编排代码。
"CompletableFuture 是 JDK 8 的异步编排框架。supplyAsync 创建有返回值的异步任务,thenApply 做同步转换(类比 map),thenCompose 做异步扁平转换(类比 flatMap),thenCombine 合并两个 Future,allOf 等全部完成,anyOf 抢最快。异常处理用 exceptionally 兜底降级、handle 双通道处理、whenComplete 做日志副作用。JDK 9 增加了 orTimeout 和 completeOnTimeout 做超时控制。生产环境三个核心原则:必须传自定义线程池、每个异步调用必须设超时、非核心模块做降级而非重试。"
Comments · 评论