CompletableFuture 异步编排
提出问题
假设你正在开发一个订单查询接口,需要同时从三个地方拿数据:用户信息(用户服务)、订单详情(订单服务)、商品快照(商品服务),等三个都拿到后,再组合成一个统一的订单视图返回。
用传统的 Future 实现的话,代码大概长这样:
ExecutorService executor = Executors.newFixedThreadPool(3);
Future<User> userFuture = executor.submit(() -> userService.getUser(uid));
Future<Order> orderFuture = executor.submit(() -> orderService.getOrder(orderId));
Future<Product> productFuture = executor.submit(() -> productService.getProduct(pid));
User user = userFuture.get(); // 阻塞
Order order = orderFuture.get(); // 阻塞
Product product = productFuture.get(); // 阻塞
return new OrderVO(user, order, product);三个 get() 都是阻塞调用。虽然三个任务确实并行执行了,但主线程在等待期间干不了别的。更糟的是,如果你想在拿到用户信息后立即查优惠券(时序依赖),Future 完全无能为力——你必须手动 get() 用户信息后再提交新任务,代码变成一层又一层的回调嵌套。
数据量级:一个典型订单查询接口,用户服务 RPC 耗时 30-80ms,订单服务 20-50ms,商品快照 10-30ms,优惠券查询 15-40ms。如果串行执行,P99 总耗时约 150-200ms;如果并行 + 异步编排,P99 可降至 60-80ms,降低 60% 以上。对于日调用量 100 万次的接口,这相当于每天节省 20-30 小时的 CPU 等待时间。
JDK 8 的 CompletableFuture 就是为了解决这类问题而生的。它让异步任务可以像 Stream API 一样链式编排,用声明式的方式描述"先做 A、A 完了做 B、同时做 C 和 D、等 B 和 C 都完成再合并"这种复杂的异步流程,而不需要显式地管理线程、锁或阻塞等待。
对比:Future vs CompletableFuture
| 维度 | Future | CompletableFuture |
|---|---|---|
| 获取结果 | 阻塞 get(),超时 get(timeout, unit) | 非阻塞 join(),链式回调,不阻塞调用线程 |
| 任务编排 | 不支持,必须手动 get() 后提交新任务 | thenApply / thenCompose / thenCombine 等声明式编排 |
| 异常处理 | 通过 ExecutionException 捕获,异常链不清晰 | exceptionally / handle / whenComplete,异常沿链传播 |
| 多任务组合 | 不支持,需手动 CountDownLatch + get() | allOf / anyOf 原生支持 |
| 超时控制 | 仅 get(timeout) 抛出 TimeoutException | JDK 9+ orTimeout / completeOnTimeout 声明式控制 |
| 主动完成 | 不支持,只能等任务结束 | complete() / completeExceptionally() / obtrudeValue() 可以手动完成 |
| 取消语义 | cancel(true/false),中断 vs 不中断 | cancel(true) 等价于 completeExceptionally(new CancellationException()) |
| 底层实现 | 基于 AbstractQueuedSynchronizer(AQS) | 基于 volatile result + CAS 无锁栈 |
分析问题
异步任务创建:supplyAsync vs runAsync
// 有返回值
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
return "Hello";
});
// 无返回值
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
System.out.println("Hello");
});两个方法都有一个带 Executor 参数的重载版本。如果不传 Executor,默认使用 ForkJoinPool.commonPool()。
这是一个常见的性能陷阱——commonPool 的线程数等于 Runtime.getRuntime().availableProcessors() - 1。如果异步任务涉及 IO(数据库查询、RPC 调用、HTTP 请求),IO 等待会占用 commonPool 的工作线程,导致其他并行流和 CompletableFuture 任务被阻塞,甚至引发线程饥饿死锁。
生产案例:某团队在 8 核容器上部署订单服务,所有 supplyAsync 都未指定线程池。高峰时 7 个 commonPool 线程全被 RPC 等待占满,导致 parallelStream 的并行流任务无法执行,整条链路超时从 50ms 飙升到 2s。修复方式就是引入自定义 IO 线程池:
// 自定义线程池,隔离 IO 和 CPU 型任务
ExecutorService ioPool = Executors.newFixedThreadPool(20);
CompletableFuture.supplyAsync(() -> userService.getUser(uid), ioPool);线程池大小估算:IO 密集型场景,线程数 = 2 × CPU 核心数 × (1 + 等待时间 / 计算时间)。以订单查询为例,RPC 等待约 60ms,计算约 5ms,比例 12:1,8 核机器建议线程数 = 2 × 8 × (1 + 12) = 208。但实际还需考虑下游服务容量,一般控制在 20-50 之间。
链式回调:thenApply / thenAccept / thenRun
拿到一个异步结果后,最常见的需求是"对它做转换"或"消费它":
CompletableFuture.supplyAsync(() -> getUser(uid))
.thenApply(user -> user.getAddress()) // 转换:User → Address
.thenAccept(address -> saveToCache(address)) // 消费:Address → void
.thenRun(() -> log("cache updated")) // 执行:不依赖前序结果
.join(); // 等待整个链完成| 方法 | 函数签名 | 返回值 | 依赖前序结果 |
|---|---|---|---|
thenApply | Function<T, R> | CompletableFuture<R> | 是 |
thenAccept | Consumer<T> | CompletableFuture<Void> | 是 |
thenRun | Runnable | CompletableFuture<Void> | 否 |
执行线程的陷阱:thenApply 默认在完成当前任务的线程上执行回调——可能是 ForkJoinPool 的 worker,也可能是主动调用 complete() 的那个线程。如果你希望回调在另一个线程池中执行,使用 thenApplyAsync 会提交到 ForkJoinPool 重新调度。thenApplyAsync(executor) 则允许指定线程池。
// thenApply → 在 supplyAsync 的线程上执行
// thenApplyAsync → 在 ForkJoinPool.commonPool 上执行
// thenApplyAsync(executor) → 在指定线程池上执行多任务组合
组合是 CompletableFuture 最强大的能力:
thenCompose:扁平化时序组合
// 不优雅:CompletableFuture<CompletableFuture<Order>>
CompletableFuture<CompletableFuture<Order>> bad =
getUser(uid).thenApply(user -> getOrder(user.getOrderId()));
// 优雅:扁平化
CompletableFuture<Order> good =
getUser(uid).thenCompose(user -> getOrder(user.getOrderId()));thenCompose 相当于 flatMap——它把返回 CompletableFuture 的回调拍平,避免嵌套。底层通过 uniComposeStage 实现,本质是注册一个 UniCompletion 节点在当前 CF 的栈上,前序完成后自动触发。
thenCombine:两个任务并行,结果合并
CompletableFuture<String> userFuture = getUser(uid);
CompletableFuture<String> orderFuture = getOrder(orderId);
CompletableFuture<OrderVO> result = userFuture
.thenCombine(orderFuture, (user, order) -> new OrderVO(user, order));thenCombine 等待两个 CompletableFuture 都完成,然后用 BiFunction 合并结果。两个任务是并行执行的(假设它们各自使用了不同的线程)。
时序图:
时间线 →
用户查询: ├─────────────────────┤
订单查询: ├───────────┤
优惠券查询: ├──────┤
(用户查完才开始)
合 并: └── 三路都完成 → 触发组装allOf:等待全部完成
CompletableFuture<User> userFuture = getUser(uid);
CompletableFuture<Order> orderFuture = getOrder(orderId);
CompletableFuture<Product> productFuture = getProduct(pid);
CompletableFuture<Void> allDone = CompletableFuture.allOf(
userFuture, orderFuture, productFuture
);
// allOf 返回 Void,需要手动获取结果
allDone.thenRun(() -> {
User user = userFuture.join(); // 此时不会阻塞,因为已经完成
Order order = orderFuture.join();
Product product = productFuture.join();
return new OrderVO(user, order, product);
});allOf 的底层实现:内部创建一个 AllOf 内部类实例,构造一个 Cnt(计数器),初始值为传入的 CF 数量。每个 CF 完成后,通过 CAS 递减计数器,当计数器归零时,调用 complete() 完成结果 CF。整个过程无锁,靠 volatile 保证可见性。
join() vs get():join() 抛出非受检的 CompletionException,get() 抛出受检的 InterruptedException 和 ExecutionException。在链式编程中 join() 更简洁,但如果异常需要传播给调用方处理,get() 的受检异常更合适。
anyOf:任意一个完成
// 多个缓存查询,哪个先返回就用哪个
CompletableFuture<Object> first = CompletableFuture.anyOf(
queryCache("redis"),
queryCache("local"),
queryCache("memcached")
);异常处理
CompletableFuture 的异常处理有三种方式,行为差异微妙:
// 1. exceptionally:只在异常时执行,提供降级值
CompletableFuture.supplyAsync(() -> riskyCall())
.exceptionally(e -> {
log.error("调用失败", e);
return "fallback";
});
// 2. handle:无论成功失败都执行,自行决定返回值
CompletableFuture.supplyAsync(() -> riskyCall())
.handle((result, e) -> {
if (e != null) {
log.error("失败", e);
return "fallback";
}
return result;
});
// 3. whenComplete:无论成功失败都执行,但不改变结果
CompletableFuture.supplyAsync(() -> riskyCall())
.whenComplete((result, e) -> {
if (e != null) {
log.error("失败", e);
} else {
log.info("成功: {}", result);
}
})
.join(); // 异常仍然会抛出| 方法 | 异常时执行 | 成功时执行 | 是否改变结果 |
|---|---|---|---|
exceptionally | 是 | 否 | 是(返回降级值) |
handle | 是 | 是 | 是(返回任意值) |
whenComplete | 是 | 是 | 否(结果不变) |
生产环境最容易踩的坑是异常静默丢失:如果 thenApply 里抛了异常,后续的 thenApply 不会执行,但异常也不会被打印——除非链上有一个 exceptionally 或 handle。调试时可以通过 get() 捕获 ExecutionException 来查看异常,但生产代码中,建议每个异步链的末端都挂一个 exceptionally 兜底。
进阶:completeExceptionally 与 obtrudeValue 的区别
completeExceptionally(Throwable ex) 通过 CAS 将 result 设为 AltResult(ex),只有任务尚未完成时才生效。obtrudeValue(T value) 则强制设置结果,即使任务已经完成也会覆盖。obtrudeValue 是破坏性的,一般只用于测试或异常恢复场景。
超时控制(JDK 9+)
// orTimeout:超时后抛出 TimeoutException
CompletableFuture.supplyAsync(() -> slowCall())
.orTimeout(3, TimeUnit.SECONDS)
.exceptionally(e -> "timeout fallback");
// completeOnTimeout:超时后使用默认值完成
CompletableFuture.supplyAsync(() -> slowCall())
.completeOnTimeout("default", 3, TimeUnit.SECONDS);JDK 8 没有原生超时方法,需要手动用 get(timeout) 或 FutureTask 包装,JDK 9 的 orTimeout 和 completeOnTimeout 是更优雅的解法。orTimeout 底层使用 Delayer(一个单线程 ScheduledExecutorService)来延迟触发 completeExceptionally(new TimeoutException())。
JDK 8 的超时替代方案:
// JDK 8 兼容方式:手动用 ScheduledExecutorService 实现超时
ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
CompletableFuture<String> timeoutFuture = future.applyToEither(
CompletableFuture.supplyAsync(() -> {
try { Thread.sleep(3000); } catch (Exception e) {}
return "timeout";
}, scheduler),
Function.identity()
);底层实现:Completion 栈
CompletableFuture 的底层是一个无锁的依赖链机制。每个 CompletableFuture 内部维护一个 volatile Object result(结果值或 AltResult 包装的异常)和一个 Completion 栈(等待当前任务完成后的回调链表)。
关键字段:
volatile Object result; // null(未完成) | T(正常结果) | AltResult(异常)
volatile Completion stack; // 回调链表栈顶(CAS 操作)完整的执行流程:
- 创建阶段:
supplyAsync(runnable)创建CompletableFuture实例,提交任务到ForkJoinPool.commonPool()(或自定义线程池) - 注册阶段:
future.thenApply(fn)创建UniApply<T,R>继承自Completion,通过CAS将自身压入前序 CF 的stack字段(栈顶) - 完成阶段:
supplyAsync的任务执行完毕,调用complete(T value):- 通过
UNSAFE.compareAndSwapObject将result从null设为T - 调用
postComplete()遍历stack链表 - 对每个
Completion节点调用tryFire(),返回true则继续传播
- 通过
- 传播阶段:
tryFire()执行回调函数,更新当前 CF 的result,然后继续调用postComplete()触发后续回调
LIFO 顺序:stack 是链表结构,push() 使用 CAS 插入栈顶,因此回调的顺序是后进先出——最后一个注册的 thenApply 最先执行。这是 CompletableFuture 的一个微妙特性。
CAS 无锁的好处:相比 FutureTask 基于 AQS 的阻塞队列,CompletableFuture 的 Completion 栈在单次完成场景下完全不需要锁,避免了线程挂起/唤醒的开销。在高吞吐场景下,这种无锁设计可以显著减少上下文切换。
实战:订单查询接口(完整版)
把开头的问题用 CompletableFuture 重写,加入超时控制、异常兜底、自定义线程池:
public OrderVO getOrderDetail(String uid, String orderId, String pid) {
ExecutorService ioPool = Executors.newFixedThreadPool(10);
CompletableFuture<User> userFuture = CompletableFuture
.supplyAsync(() -> userService.getUser(uid), ioPool);
CompletableFuture<Order> orderFuture = CompletableFuture
.supplyAsync(() -> orderService.getOrder(orderId), ioPool);
CompletableFuture<Product> productFuture = CompletableFuture
.supplyAsync(() -> productService.getProduct(pid), ioPool);
// 用户信息查完后,立即查优惠券(时序依赖)
CompletableFuture<Coupon> couponFuture = userFuture
.thenApplyAsync(user -> couponService.getCoupon(user.getLevel()), ioPool);
// 等待三个主要数据 + 优惠券全部就绪,组装结果
CompletableFuture<OrderVO> result = CompletableFuture
.allOf(userFuture, orderFuture, productFuture, couponFuture)
.thenApplyAsync(v -> {
User user = userFuture.join();
Order order = orderFuture.join();
Product product = productFuture.join();
Coupon coupon = couponFuture.join();
return new OrderVO(user, order, product, coupon);
}, ioPool)
.exceptionally(e -> {
log.error("组装订单视图失败", e);
return OrderVO.empty(); // 降级返回空对象
})
.orTimeout(5, TimeUnit.SECONDS)
.exceptionally(e -> {
log.warn("订单查询超时,uid={}, orderId={}", uid, orderId);
return OrderVO.empty();
});
return result.join();
}这个设计的优势:
- 三个无依赖的查询并行执行,总耗时 = max(用户服务, 订单服务, 商品服务) ≈ 60ms
- 用户信息一拿到,自动触发优惠券查询(不阻塞,时序依赖自动处理)
- 所有数据就绪后自动组装(不主动轮询)
- 异常和超时都有兜底,不会让调用方空等
ioPool隔离了 IO 线程,不影响 commonPool 的并行流计算
更现实的线程池隔离策略:
// 按下游服务隔离线程池,避免一个服务雪崩拖垮其他服务
ExecutorService userPool = Executors.newFixedThreadPool(5, r -> new Thread(r, "user-rpc-"));
ExecutorService orderPool = Executors.newFixedThreadPool(5, r -> new Thread(r, "order-rpc-"));
ExecutorService productPool = Executors.newFixedThreadPool(3, r -> new Thread(r, "product-rpc-"));高级用法
生产级重试模式
public <T> CompletableFuture<T> withRetry(
Supplier<T> action, int maxRetries, Executor executor) {
CompletableFuture<T> cf = CompletableFuture.supplyAsync(action, executor);
for (int i = 0; i < maxRetries; i++) {
cf = cf.exceptionally(e -> {
log.warn("重试 {}/{}", i + 1, maxRetries, e);
return CompletableFuture.supplyAsync(action, executor).join();
});
}
return cf;
}熔断式超时(JDK 9+)
// 先快速返回缓存,后台异步刷新
CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> fetchFromDb())
.completeOnTimeout(cache.get(key), 100, TimeUnit.MILLISECONDS)
.thenApply(result -> {
cache.put(key, result); // 后台刷新缓存
return result;
});总结
关键点清单
supplyAsync/runAsync创建异步任务,务必指定自定义线程池处理 IO 型任务,避免公共池被阻塞thenApply(转换)、thenAccept(消费)、thenRun(执行)构建链式回调;Async后缀的版本切换到新线程执行thenCompose扁平化时序组合,避免CompletableFuture<CompletableFuture>嵌套thenCombine并行执行两个任务后合并结果;allOf等待全部完成;anyOf等待任意一个完成- 异常处理三件套:
exceptionally(异常降级)、handle(成功失败都处理)、whenComplete(感知但不改变结果) - 异步链末端必须挂异常处理,否则异常会静默丢失
- JDK 9+ 的
orTimeout和completeOnTimeout提供原生超时控制 - 底层基于
volatile result+ CAS 无锁传播的Completion栈
面试话术示例
"CompletableFuture 的核心价值是把异步编排从回调地狱变成声明式链式调用。底层的
Completion栈是无锁的,通过 CAS 更新volatile result触发后续回调传播。生产中最关键的是两件事:一是 IO 密集型任务必须用自定义线程池,别用commonPool把所有并行流都堵住——我见过 8 核机器上 commonPool 7 个线程全被 RPC 等待占满导致整条链路超时翻 40 倍的案例;二是每个异步链末端必须放一个exceptionally兜底,否则异常会丢得无声无息。JDK 9 的orTimeout是超时控制的最佳实践,JDK 8 只能靠get(timeout)或ScheduledExecutorService手动模拟。allOf底层通过 CAS 计数器实现无锁完成检测,join()比get()在链式编程中更简洁因为抛的是非受检异常。"
参考:
CompletableFuture源码java.util.concurrent.CompletableFuture、CompletionStage接口文档、ForkJoinPool.commonPool()源码、JDK 9 JEP 266(CompletableFuture 增强)