首页 / Java 学习笔记 / 14

JAVA · Vol.II · DAY 14 · 并发编程

CompletableFuture 异步编排:从回调地狱到流式组合

中级高频#并发#异步
第 1 站

开场:商品详情页为什么要 800ms?

某次大促压测,商品详情页接口 P99 高达 1.2 秒,是竞品的三倍。架构师打开链路监控一看:四路后端调用是串行跑的。为什么不能一起发出去?

一个商品详情页,页面渲染需要 4 个后端数据,它们之间完全没有依赖关系

数据模块调用方式平均耗时
商品基础信息DB 查询200ms
价格与优惠价格中心 RPC200ms
用户评价评价服务 RPC200ms
推荐商品推荐引擎 RPC200ms

如果按顺序串行调用,总耗时 = 4 × 200ms = 800ms;而这 4 个调用彼此独立,完全可以在同一个时刻一起发出,等全部完成后再组装页面——这就是典型的「扇出 / 扇入(Fan-out / Fan-in)」。CompletableFuture 正是为这类「异步编排」而生的。

串行 vs 并行——同一个接口的两种写法
// 串行: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 到底差在哪。

第 2 站

传统 Future 的三大痛点:阻塞、无回调、不能组合

JDK 5 引入的 Future 已经能做到「任务丢到线程池、之后取结果」,但用它写「多个异步任务的编排」非常别扭。先看代码,再总结痛点:

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,代码里全是噪音。

第 3 站

回调地狱:嵌套回调为什么让人抓狂

既然 Future 没有回调,当年没有 CompletableFuture 的开发者是怎么解决「结果就绪自动处理」的?答案是回调(Callback)——把「下一步做什么」作为参数传给异步方法。但回调一旦嵌套,就成了著名的回调地狱(Callback Hell)

回调嵌套示例 · 三个有依赖的异步步骤(以 Guava ListenableFuture 风格示意)
// 查用户 → 查会员价 → 查优惠券,每步都依赖上一步
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 站,同一段需求两种写法摆在一起看,你会立刻明白它解决了什么。

第 4 站

CompletableFuture 是什么:Future + 完成即触发

CompletableFuture 同时实现了两个接口:Future(能拿结果、能取消)+ CompletionStage(完成阶段:能在「完成」这个事件上挂一串后续动作)。它的核心机制一句话——完成即触发(Completion):任务完成的那一刻,所有提前登记好的依赖动作会被自动逐个触发,不需要你手动通知、不需要轮询。

任务线程 跑完任务 / 或外部线程 f.complete(结果) 写结果 + 唤醒所有回调 CompletableFuture volatile Object result (结果 或 异常) 依赖动作栈 (无锁 Treiber 栈) 回调 1:thenAccept 消费结果 回调 2:thenApply 转换结果 回调 3:exceptionally 捕获异常 谁先 complete,谁负责唤醒全部登记过的回调 —— 无需轮询、无需手动通知
图 1CompletableFuture 的「完成即触发」模型:结果是字段,回调是栈,完成时统一唤醒

内部实现上(JDK 8),它就是两个字段:一个 volatile Object result 存结果(或异常),一个 volatile Completion stack 存所有依赖动作——用无锁的 Treiber 栈组织,谁调 complete(),谁负责把这个栈弹出来逐个执行。所以「注册回调」和「完成任务」完全解耦,可以发生在任何线程、任何顺序:

手动完成示例 · 先注册回调,后 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()
线程池控制提交时指定创建 + 每个编排动作都能指定
第 5 站

创建 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 短路,下游完全不用感知「这次是缓存还是真的查了」:

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);
}
方式入参返回是否异步执行典型场景
supplyAsyncSupplier<T>CompletableFuture<T>有返回值的异步任务(查 DB、调 RPC)
runAsyncRunnableCompletableFuture<Void>无返回值的异步动作(发通知、写日志)
completedFutureTCompletableFuture<T>否(立即完成)测试替身、缓存短路、默认值
new + complete()手动控制完成时机(配合定时器做超时)
核心记忆

supplyAsync 有返回值,runAsync 无返回值,completedFuture 立即可用。补充:两个 Async 创建方法都有带 Executor 的重载——生产环境必须用这个重载,原因下一站讲。

第 6 站

线程池选择:为什么生产必须传自定义线程池

supplyAsync / runAsync 不传 Executor 时,默认使用 ForkJoinPool.commonPool()——这是全 JVM 共享的一个公共线程池,线程数 = CPU 核数 - 1(可用 -Djava.util.concurrent.ForkJoinPool.common.parallelism 覆盖)。听起来很方便,生产里却是个大坑:

commonPool 的危险用法 · IO 密集任务挤爆公共池
// 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 这类回调到底跑在哪条线程上?

点破:完成线程 vs 线程池

不带 Async 后缀的方法(thenApply / thenAccept 等)在「触发完成的那条线程」上同步执行——谁 complete 的,谁顺手把回调也跑了;带 Async 后缀的方法(thenApplyAsync 等)才丢进线程池。所以:回调里如果要做耗时操作,要么用 Async 变体、要么保证完成线程不阻塞(比如别在 RPC 的 IO 线程里做重活)。这也是很多线上 bug 的来源:回调在 IO 线程里跑了 500ms 的 CPU 计算,把下游连接池拖垮。

核心记忆

生产环境必须传自定义 Executor,不要用默认的 ForkJoinPool.commonPool()——线程数不可控、全 JVM 共享、互相干扰。Async 变体换线程池执行,非 Async 变体在完成线程上执行。

第 7 站

链式转换:thenApply / thenAccept / thenRun

CompletableFuture 最强大的特性是链式组合——像 Stream 一样把多个异步步骤串成流水线(Pipeline)。面试最爱问的三个方法:thenApply(转换)、thenAccept(消费)、thenRun(执行动作)。

thenApply / thenApplyAsync / 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"));
thenApplymap(T → U,有返回值) |  thenAcceptforEach(T → void,消费) |  thenRun(跟结果无关,跑个动作)

选型口诀:要转换结果用 thenApply,要消费结果用 thenAccept,跟结果没关系用 thenRun。方法名后缀规律也要记住:Async 表示丢线程池、不带则同步执行(执行线程的语义见第 6 站)。

核心记忆

thenApply 有参有返回、thenAccept 有参无返回、thenRun 无参无返回。它们都只适合「下一步是同步操作」——如果下一步本身是异步的(返回 CompletableFuture),就要用下一站的 thenCompose。

第 8 站

thenCompose vs thenApply:嵌套与拍平

这是面试出现频率最高的追问:「thenApply 和 thenCompose 有什么区别?」一句话版:thenApply 的回调返回普通值 U,thenCompose 的回调返回 CompletableFuture<U>——thenCompose 会把返回的 Future 自动「拍平」,避免嵌套。

thenApply vs thenCompose——最关键的区别
// 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 接异步操作
// ✗ 错误:用 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();   // 一层搞定
类比:Stream 的 map vs flatMap

thenCompose 就是异步版的 flatMapflatMap 把「元素 → 集合」拍平成一层流,thenCompose 把「值 → Future」拍平成一层 Future。判断口诀:回调里是同步方法用 thenApply,回调里返回 Future 用 thenCompose——永远不要让链上出现 CompletableFuture<CompletableFuture<…>>

核心记忆

thenApply 回调返回普通值(map),thenCompose 回调返回 Future 并自动拍平(flatMap)。下一步是异步操作用 thenCompose,否则会得到嵌套的 CompletableFuture,取结果被迫二次阻塞。

第 9 站

组合与聚合:thenCombine / allOf / anyOf

链式转换解决的是「串行流水线」,但真实场景中经常需要多个独立 Future 的聚合。这就是 thenCombineallOfanyOf 的舞台。

thenCombine Future A Future B A + B → C allOf F1 F2 F3 F4 全部完成 anyOf F1 F2 F3 F4 最先完成 两个 Future 合并为一个结果 等待所有 Future 完成 任一 Future 完成即返回 实战:商品 + 价格 → 商品详情 VO Product Future Price Future thenCombine → VO
图 2CompletableFuture 三种组合模式——thenCombine / allOf / anyOf
thenCombine:两个 Future 合并为一个结果
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;
});
allOf:等待全部完成
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
anyOf:最先完成的胜出
// 多数据源竞速:谁先返回用谁(常用于缓存 + 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 / anyOf 的两个细节

allOf 的失败语义:只要其中一个 Future 异常完成,返回的 all 也异常完成——所以 all.thenRun(...) 里的回调不会执行;但其它已经成功完成的任务不受影响,各自的 join() 仍然能取到结果(谁失败谁抛)。② anyOf 的应用:不止「缓存 vs DB」双读,多机房容灾「取最快到达的那路」也是它——但注意要配合超时,否则最慢那路永远不完成,anyOf 本身不设超时也会一直等。

核心记忆

thenCombine 是「两路合并」,allOf 是「等所有」,anyOf 是「抢最快」。allOf 返回 Void 需要手动取结果,anyOf 返回 Object 需要强转。

第 10 站

异常处理三兄弟: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);
    });  // 结果原样往下传,异常也原样往下传
方法输入输出是否改变结果典型场景
exceptionallyThrowable同类型 T是(兜底值)降级返回默认值
handle(T, Throwable)新类型 U成功/失败都需要转换
whenComplete(T, Throwable)同类型 T日志、埋点、监控
常见踩坑:exceptionally 的返回类型

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 才爆发。

第 11 站

get() vs join():受检异常与运行时异常

get()join() 都能「阻塞地取结果」,区别全在异常处理方式上,这也是面试爱问的细节:

get() vs 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 + ExecutionExceptionCompletionException
是否受检受检(checked),必须 try-catch非受检(unchecked),可不处理
中断语义可中断:等待中被 interrupt 会抛异常不可中断:中断只被记下,等完再恢复标志
异常包装原始异常在 ExecutionException.getCause()原始异常在 CompletionException.getCause()
典型场景最外层 Controller 取最终结果,显式处理异常链式内部、测试代码,追求简洁

还有一个不阻塞的变体值得记住——getNow(valueIfAbsent):Future 已完成就返回结果,未完成就直接返回你给的默认值,绝不阻塞。适合「尽力而为」的读取:

getNow:不阻塞的"尽力而为"读取
// 缓存异步预热中:读到了就用,没读到先返回默认值,绝不阻塞请求线程
CompletableFuture<Price> preheating = loadPrice(id);          // 可能还没完成
Price price = preheating.getNow(Price.UNKNOWN);                // 未完成 → 立即返回 UNKNOWN
核心记忆

get() 抛受检异常(InterruptedException + ExecutionException),join() 抛非受检的 CompletionException;两者都把原始异常包在 getCause() 里。链路内部建议统一 join(),最外层用 get() 显式兜底;需要「不阻塞拿结果」用 getNow()。

第 12 站

超时控制:orTimeout / completeOnTimeout 与 JDK 8 方案

生产环境中,异步调用必须有超时——否则一个下游服务 hang 住,你的线程会被无限阻塞,最终线程池耗尽、服务雪崩。没有超时的异步调用,就像挂机等一个永远不会来的电话。

JDK 9+:orTimeout 与 completeOnTimeout
// 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 兼容方案:ScheduledExecutorService
// 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 手动实现。生产环境每个异步调用都必须设超时——这是防止线程池耗尽雪崩的底线。

第 13 站

生产实战:商品详情页完整编排

面试中如果能写出以下完整模式,会非常有加分。我们用商品详情页把三种核心编排模式串起来(以下模板假设 JDK 9+;JDK 8 请把 orTimeout / completeOnTimeout 换成第 12 站的 withTimeout 工具方法):

完整实战:商品详情页编排(Pipeline + Fan-out/Fan-in + 超时降级)
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;
        });
}
模式 3:Timeout + Retry(超时重试)
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);
模式特点适用场景
PipelinethenCompose 串联,前一步输出是后一步输入有依赖的多步异步(如:查用户 → 查会员价)
Fan-out / Fan-inallOf 并行发起,等全部完成再聚合无依赖的多路聚合(如:商品 + 评价 + 推荐)
Timeout + RetryorTimeout + exceptionally 递归重试不稳定的下游调用,需要容错
面试官追问:如果推荐服务挂了怎么办?

推荐是非核心模块。用 completeOnTimeout 超时降级为空列表 + exceptionally 捕获异常也返回空列表。核心模块(商品、价格)才需要重试。原则:核心模块重试 + 超时,非核心模块降级 + 超时——重试也要设上限,且最好退避,否则雪崩时你的重试就是在补刀。

生产坑:traceId 串号

CF 的回调经常在别的线程上执行(线程池线程、完成线程),而日志框架的 traceId 通常放在 ThreadLocal 里——跨线程就丢了,全链路追踪在异步段断链,排查问题抓瞎。解决:① 提交任务前把 traceId 抓出来,回调里手动塞回;② 用阿里开源的 TransmittableThreadLocal 做线程池透传(卷二 ThreadLocal 篇有专题)。

第 14 站

回调地狱 vs 流式编排:同一个需求两种写法

把第 3 站那个「查用户 → 查会员价 → 算实付」的需求,用两种方式各写一遍。先看对比,再谈感受:

写法 A:回调嵌套 —— 三层缩进,错误处理重复
// 每一步都要:传回调 + 传错误处理,越嵌越深
fetchUser(userId, user ->
    fetchMemberPrice(user, price ->
        calcPay(price, pay -> render(pay), onError),
    onError),
onError);
写法 B:CompletableFuture 流式编排 —— 一维链,异常集中处理
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)」的全部意义:不是少了几个括号,而是让流程可读、可复用、可组合。

第 15 站

方法速查表:一张表记住全部 API

面试前把这张表过一遍,主要 API 就不会漏:

CompletableFuture 方法速查
// 还没列全的"双 Future 变体":thenAcceptBoth(两个都完成再消费)、
// applyToEither / acceptEither(两个任一完成就处理)。面试提到"还有变体"即可加分。
类别方法一句话说明Stream 类比
创建supplyAsync有返回值的异步任务
runAsync无返回值的异步任务
completedFuture已完成的 Future(测试/默认值/短路)Stream.of()
链式转换thenApply同步转换 T → Umap
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。

第 16 站

常见坑清单与最佳实践

常见坑清单
  1. 忘记传自定义线程池——默认 ForkJoinPool.commonPool() 线程数 = CPU 核数 - 1,IO 密集场景必炸,还拖垮别人的 parallelStream
  2. allOf 返回 Void——不能直接从 allOf 拿结果,必须手动 join() 每个子 Future
  3. exceptionally 返回 null——下游收到 null 引发 NPE,应返回有意义的默认值
  4. 没有设超时——下游 hang 住导致线程永久阻塞,最终线程池耗尽、服务雪崩
  5. thenApply vs thenCompose 混用——下一步是异步操作用 thenCompose,否则得到嵌套 Future,取结果二次阻塞
  6. join() / get() 选错——get() 是受检异常必须处理;链路内部统一 join() 更干净
  7. 在非 Async 回调里做耗时操作——回调在完成线程上跑,会阻塞那条线程(可能是 RPC 的 IO 线程)
  8. ThreadLocal 不传递——回调跨线程,traceId 丢失、全链路追踪断链(用 TransmittableThreadLocal)
  9. 异常被吞——没接住也没 get/join 的异常静默沉底,线上只剩"结果不对"的谜案;记得打日志或上报
  10. 重试无上限 / 无退避——雪崩时疯狂重试等于补刀,重试必须设上限 + 退避 + 熔断
最佳实践:把超时 + 重试封装成统一工具
// 团队内统一封装,禁止裸写 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 更低、可读性更高、可维护的异步编排代码。

面试 30 秒总结

"CompletableFuture 是 JDK 8 的异步编排框架。supplyAsync 创建有返回值的异步任务,thenApply 做同步转换(类比 map),thenCompose 做异步扁平转换(类比 flatMap),thenCombine 合并两个 Future,allOf 等全部完成,anyOf 抢最快。异常处理用 exceptionally 兜底降级、handle 双通道处理、whenComplete 做日志副作用。JDK 9 增加了 orTimeout 和 completeOnTimeout 做超时控制。生产环境三个核心原则:必须传自定义线程池、每个异步调用必须设超时、非核心模块做降级而非重试。"

Comments · 评论