CompletableFuture 全过程原理与商业异步编排
CompletableFuture 是 JDK 8 引入的异步编排工具。很多人会用 supplyAsync、thenApply、allOf,但真正上项目时更重要的是说清:
- 它和
Future的本质区别是什么。 - 异步任务到底由哪个线程池执行。
thenApply、thenApplyAsync的线程有什么区别。- 多个任务怎么组合,依赖任务怎么串联。
- 异常怎么传播,为什么
allOf后还要逐个取结果。 - 超时是不是会真正取消底层任务。
- 为什么不能把阻塞 IO 都丢到公共
ForkJoinPool。 - 接口聚合、批量采集、异步入库这些商业场景怎么用才安全。
一句话先建立主线:
CompletableFuture 不是新线程,也不是性能魔法。它是“异步任务结果 + 依赖回调链”的编排对象,真正执行任务的仍然是线程池。
学习目标
学完这一页,你要能回答:
Future为什么不适合复杂异步编排。CompletableFuture内部保存哪些状态。supplyAsync、runAsync、thenApply、thenCompose、thenCombine、allOf分别解决什么问题。- 非 Async 回调和 Async 回调由哪个线程执行。
join()和get()有什么区别。exceptionally、handle、whenComplete有什么区别。orTimeout、completeOnTimeout的真实含义。- 如何写一个可落地的订单详情接口聚合 Demo。
- 线上异步任务堆积、线程池打满、下游慢、异常被吞怎么排查。
Future 的痛点
JDK 5 就有 Future,但它更像“任务回执单”。
Future<String> userFuture = executor.submit(() -> queryUser());
String user = userFuture.get();它能拿结果,但有几个问题:
| 痛点 | 说明 | 后果 |
|---|---|---|
get() 阻塞 | 调用方线程要等结果 | 容易把线程池卡住 |
| 不擅长组合 | 多个 Future 要手动逐个 get | 代码复杂,异常难处理 |
| 不擅长依赖 | 任务 B 依赖任务 A 的结果要手写回调 | 容易写成阻塞等待 |
| 不擅长异常链 | 异常处理分散在多个 get 周围 | 兜底逻辑混乱 |
| 不能主动完成 | 只能由任务执行完成 | 不方便超时兜底和外部补偿 |
CompletableFuture 的核心增强是:它既表示未来结果,又能描述后续依赖动作。
flowchart TD
A["Future<br/>只能等结果"] --> B["get 阻塞"]
C["CompletableFuture"] --> D["保存结果或异常"]
C --> E["注册后续动作"]
C --> F["组合多个任务"]
C --> G["主动 complete"]CompletableFuture 内部可以怎么理解
从学习角度,可以把一个 CompletableFuture 理解成四部分:
flowchart TD
A["CompletableFuture"] --> B["任务状态<br/>未完成/成功/失败/取消"]
A --> C["结果或异常<br/>value / Throwable"]
A --> D["依赖动作栈<br/>thenApply / thenAccept 等"]
A --> E["执行器<br/>默认公共池或自定义线程池"]任务完成时,它会:
- 把状态改成成功或失败。
- 保存结果或异常。
- 触发依赖它的后续动作。
- 把结果或异常继续向后传播。
流程:
flowchart TD
A["创建 CompletableFuture"] --> B["注册依赖动作"]
B --> C["异步任务在线程池执行"]
C --> D{"任务成功还是失败"}
D -- "成功" --> E["保存正常结果"]
D -- "失败" --> F["保存异常结果"]
E --> G["触发后续阶段"]
F --> H["触发异常处理阶段"]
G --> I["继续向后传播"]
H --> IsupplyAsync 和 runAsync
| 方法 | 是否有返回值 | 适合场景 |
|---|---|---|
supplyAsync | 有 | 查询用户、订单、库存并返回结果 |
runAsync | 无 | 异步写日志、发通知、刷新缓存 |
Demo:
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class SupplyRunDemo {
private static final ExecutorService POOL = Executors.newFixedThreadPool(4);
public static void main(String[] args) {
CompletableFuture<String> userFuture =
CompletableFuture.supplyAsync(() -> "user-1001", POOL);
CompletableFuture<Void> logFuture =
CompletableFuture.runAsync(() -> System.out.println("write log"), POOL);
System.out.println(userFuture.join());
logFuture.join();
POOL.shutdown();
}
}生产里不要直接使用 Executors.newFixedThreadPool 作为最终方案,这里只是最小 Demo。真实项目要使用有界队列、自定义线程名、拒绝策略和监控。
非 Async 和 Async 的线程区别
这是 CompletableFuture 面试和线上排查最重要的点之一。
| 写法 | 后续阶段由谁执行 |
|---|---|
thenApply | 通常由完成上一步的线程继续执行 |
thenApplyAsync 不传线程池 | 默认提交到公共 ForkJoinPool |
thenApplyAsync 传线程池 | 提交到指定线程池 |
流程:
flowchart TD
A["任务 A 在线程 pool-1 执行"] --> B["任务 A 完成"]
B --> C{"后续是 thenApply 还是 thenApplyAsync"}
C -- "thenApply" --> D["可能继续由 pool-1 执行"]
C -- "thenApplyAsync 无 executor" --> E["提交到公共 ForkJoinPool"]
C -- "thenApplyAsync 有 executor" --> F["提交到指定线程池"]为什么这很重要?
- 如果
thenApply里做慢 IO,可能拖慢完成上一步的线程。 - 如果
thenApplyAsync不指定线程池,阻塞任务会占用公共池。 - 公共池被占满会影响其他并行 Stream、其他 CompletableFuture 任务。
建议:
| 场景 | 建议 |
|---|---|
| 轻量内存转换 | thenApply |
| 下游 IO 调用 | thenComposeAsync 或 thenApplyAsync + 自定义 IO 线程池 |
| CPU 计算 | 独立 CPU 线程池 |
| 核心接口聚合 | 必须自定义线程池和超时 |
thenApply、thenCompose、thenCombine 怎么选
thenApply:同步转换结果
上一步返回用户 ID,下一步只是拼接字符串:
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> 1001L)
.thenApply(userId -> "user-" + userId);thenApply 的返回值不是新的异步任务,而是普通转换结果。
thenCompose:串联依赖的异步任务
用户详情依赖用户 ID,下一步本身也是异步查询:
CompletableFuture<String> userDetail = queryUserId()
.thenCompose(userId -> queryUserDetail(userId));如果用 thenApply,会得到嵌套:
CompletableFuture<CompletableFuture<String>> wrong =
queryUserId().thenApply(userId -> queryUserDetail(userId));所以:
| 场景 | 应该用 |
|---|---|
| A 结果转成普通 B | thenApply |
| A 结果决定下一个异步 Future | thenCompose |
thenCombine:合并两个独立任务
订单和优惠券互不依赖,可以并行后合并:
CompletableFuture<String> merged = queryOrder()
.thenCombine(queryCoupon(), (order, coupon) -> order + " + " + coupon);流程图:
flowchart TD
A["查询订单"] --> C["thenCombine 合并"]
B["查询优惠券"] --> C
C --> D["生成订单详情片段"]allOf 的真实含义
allOf 表示“等这些 Future 都完成”,但它不帮你把每个结果收集成列表。
CompletableFuture<Void> all =
CompletableFuture.allOf(userFuture, orderFuture, couponFuture);
all.join();为什么 allOf 不直接返回 List<T>?
因为每个 Future 的类型可能不同:
CompletableFuture<UserDTO> userFuture;
CompletableFuture<OrderDTO> orderFuture;
CompletableFuture<List<CouponDTO>> couponFuture;它只能统一表示“都完成了”,具体结果还要从每个 Future 中取。
allOf 遇到异常
只要其中一个 Future 异常,allOf().join() 也会异常。但其他任务可能已经完成,也可能还在继续执行。
flowchart TD
A["userFuture 成功"] --> D["allOf"]
B["orderFuture 异常"] --> D
C["couponFuture 仍执行"] --> D
D --> E["allOf 以异常完成"]所以生产代码要考虑:
- 哪些任务失败可以兜底。
- 哪些任务失败必须整体失败。
- 是否要记录每个子任务的异常。
- 是否要对慢任务设置超时。
异常处理方法区别
| 方法 | 成功时调用 | 失败时调用 | 能否改变结果 |
|---|---|---|---|
exceptionally | 否 | 是 | 能,返回兜底值 |
handle | 是 | 是 | 能,返回新结果 |
whenComplete | 是 | 是 | 通常只观察,不改变上游结果 |
Demo:
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> {
throw new RuntimeException("query failed");
})
.exceptionally(ex -> "default-user");
System.out.println(future.join());whenComplete 常用于日志:
CompletableFuture<String> future = queryUser()
.whenComplete((result, error) -> {
if (error != null) {
System.out.println("query user failed: " + error.getMessage());
}
});注意:whenComplete 里如果自己又抛异常,可能覆盖或影响后续异常表现。日志逻辑要稳,不能在异常处理里再制造异常。
join 和 get 区别
| 对比 | join() | get() |
|---|---|---|
| 异常类型 | 抛 CompletionException | 抛 ExecutionException、InterruptedException |
| 是否受检异常 | 不是 | 是 |
| 中断处理 | 不强制你处理 | 必须处理 InterruptedException |
| 常见使用 | 函数式链路里更简洁 | 传统 Future 风格 |
生产代码里,如果用 get() 捕获 InterruptedException,要恢复中断标记:
try {
return future.get();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("interrupted", e);
}超时控制:不是所有超时都会杀掉底层任务
JDK 9 之后有:
| 方法 | 行为 |
|---|---|
orTimeout | 超时后让 Future 以异常完成 |
completeOnTimeout | 超时后用默认值完成 |
示例:
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> slowQuery(), IO_POOL)
.completeOnTimeout("default", 800, TimeUnit.MILLISECONDS);重点:
Future 超时完成,不一定代表底层正在执行的任务被强制停止。
如果 slowQuery() 已经在线程中执行,并且底层 HTTP/DB 调用没有超时,它可能还会继续占用线程。真正的生产超时要两层都做:
- CompletableFuture 编排层超时,避免接口一直等。
- HTTP、DB、Redis、RPC 客户端自身超时,避免工作线程长期卡住。
商业 Demo:订单详情接口聚合
这个 Demo 模拟订单详情页并行查询用户、订单、优惠券,并处理超时和异常。
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class OrderDetailCompletableFutureDemo {
private static final ThreadPoolExecutor IO_POOL = new ThreadPoolExecutor(
8,
16,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<Runnable>(200),
r -> new Thread(r, "order-aggregate-" + System.nanoTime()),
new ThreadPoolExecutor.CallerRunsPolicy()
);
public static void main(String[] args) {
OrderDetail detail = queryOrderDetail(1001L);
System.out.println(detail);
IO_POOL.shutdown();
}
private static OrderDetail queryOrderDetail(Long orderId) {
CompletableFuture<String> userFuture = CompletableFuture
.supplyAsync(() -> queryUser(orderId), IO_POOL)
.completeOnTimeout("unknown-user", 800, TimeUnit.MILLISECONDS)
.exceptionally(ex -> "unknown-user");
CompletableFuture<String> orderFuture = CompletableFuture
.supplyAsync(() -> queryOrder(orderId), IO_POOL)
.orTimeout(1, TimeUnit.SECONDS);
CompletableFuture<String> couponFuture = CompletableFuture
.supplyAsync(() -> queryCoupon(orderId), IO_POOL)
.completeOnTimeout("no-coupon", 500, TimeUnit.MILLISECONDS)
.exceptionally(ex -> "no-coupon");
CompletableFuture<Void> all = CompletableFuture.allOf(userFuture, orderFuture, couponFuture);
try {
all.join();
return new OrderDetail(
userFuture.join(),
orderFuture.join(),
couponFuture.join()
);
} catch (Exception e) {
throw new RuntimeException("query order detail failed, orderId=" + orderId, e);
}
}
private static String queryUser(Long orderId) {
sleep(100);
return "user-of-" + orderId;
}
private static String queryOrder(Long orderId) {
sleep(120);
return "order-" + orderId;
}
private static String queryCoupon(Long orderId) {
sleep(80);
return "coupon";
}
private static void sleep(long millis) {
try {
TimeUnit.MILLISECONDS.sleep(millis);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}
static class OrderDetail {
private final String user;
private final String order;
private final String coupon;
OrderDetail(String user, String order, String coupon) {
this.user = user;
this.order = order;
this.coupon = coupon;
}
public String toString() {
return "OrderDetail{user='" + user + "', order='" + order + "', coupon='" + coupon + "'}";
}
}
}这个 Demo 体现的生产原则:
| 原则 | 说明 |
|---|---|
| 自定义线程池 | 避免阻塞公共池 |
| 有界队列 | 防止无限堆积 |
| 线程命名 | 方便 jstack 排查 |
| 局部兜底 | 非核心数据可返回默认值 |
| 核心失败抛异常 | 订单主信息失败不能假装成功 |
| 编排层超时 | 防止接口无限等待 |
真实项目还要加 traceId、日志、指标、下游客户端超时、熔断限流和降级策略。
商业场景怎么选
| 场景 | 是否适合 CompletableFuture | 注意点 |
|---|---|---|
| 订单详情接口聚合 | 适合 | 下游互不依赖,并且有超时兜底 |
| 批量采集多个医院接口 | 适合但要谨慎 | 控制并发,不能打爆目标系统 |
| 支付扣款主流程 | 谨慎 | 不能随便异步化,要保证状态一致 |
| 异步写日志/埋点 | 适合 | 可降级,不能影响主流程 |
| 大量 CPU 计算 | 可用 | 使用 CPU 线程池,避免和 IO 混用 |
| MQ 消费内部并发 | 谨慎 | offset 提交、幂等、失败重试要设计好 |
线上排查
现象一:接口用了 CompletableFuture 仍然很慢
排查路径:
flowchart TD
A["接口仍然慢"] --> B["看每个子任务耗时"]
B --> C{"是否某个下游 P99 很高"}
C -- "是" --> D["查下游超时、连接池、慢 SQL"]
C -- "否" --> E["看聚合线程池 active 和 queue"]
E --> F{"线程池是否满"}
F -- "是" --> G["降并发、扩容或隔离线程池"]
F -- "否" --> H["查是否 join/get 放错位置导致串行等待"]常见错误:
String user = CompletableFuture.supplyAsync(() -> queryUser(), POOL).join();
String order = CompletableFuture.supplyAsync(() -> queryOrder(), POOL).join();这其实是提交一个等一个,仍然接近串行。正确做法是先创建所有 Future,再统一等待。
现象二:公共 ForkJoinPool 被打满
原因:
supplyAsync没传线程池。parallelStream和 CompletableFuture 共用公共池。- 任务里做阻塞 IO。
处理:
- 业务异步必须传自定义线程池。
- IO 任务和 CPU 任务分池。
- 给线程命名,接入监控。
现象三:allOf 卡住或异常
常见原因:
| 原因 | 处理 |
|---|---|
| 子任务没有超时 | 加编排层和客户端层超时 |
| 子任务异常未处理 | 根据业务选择兜底或整体失败 |
| 子任务线程池队列堆积 | 看 active、queue、拒绝次数 |
| 子任务里又等待同池任务 | 避免线程饥饿死锁 |
现象四:超时后线程仍然被占用
原因是 Future 超时完成不等于底层 IO 被强制中断。要检查:
- HTTP read timeout。
- DB query timeout。
- RPC timeout。
- 任务是否响应中断。
- 是否调用了
cancel(true),以及任务是否能响应。
常见坑
| 坑 | 后果 | 正确做法 |
|---|---|---|
| 不传线程池 | 阻塞公共池 | 自定义业务线程池 |
| 在链路中提前 join | 并行退化成串行 | 先创建 Future,再统一等待 |
| allOf 后不处理子任务异常 | 不知道哪个任务失败 | 每个子任务加日志和兜底 |
| 异步任务无超时 | 线程长期占用 | 编排层和客户端层都超时 |
| IO 和 CPU 混用一个池 | 互相拖慢 | 按任务类型隔离 |
| 兜底返回默认成功 | 掩盖核心失败 | 区分核心和非核心依赖 |
| 忽略 traceId 传递 | 日志链路断裂 | 包装任务或使用上下文传递组件 |
面试标准回答
CompletableFuture 是什么,和 Future 有什么区别
标准回答:
Future 表示一个未来结果,主要通过 get 阻塞等待结果,不擅长任务组合和依赖编排。CompletableFuture 是 JDK 8 提供的异步编排工具,既表示未来结果,也可以注册 thenApply、thenCompose、thenCombine、allOf 等后续阶段,把多个异步任务组织成依赖链或并行聚合。它本身不是新线程,真正执行任务的仍然是线程池。生产使用时必须指定业务线程池,处理异常和超时,避免阻塞公共 ForkJoinPool。thenApply、thenCompose、thenCombine 区别
标准回答:
thenApply 用于把上一步结果同步转换成另一个普通结果;thenCompose 用于上一步结果决定下一个异步任务,可以把嵌套的 CompletableFuture 拉平;thenCombine 用于两个互不依赖的异步任务都完成后合并结果。简单说,普通转换用 thenApply,依赖异步串联用 thenCompose,两个并行结果合并用 thenCombine。CompletableFuture 线上使用注意什么
标准回答:
线上使用 CompletableFuture 要注意五点:第一,必须指定自定义线程池,避免阻塞公共 ForkJoinPool;第二,先创建多个 Future 再统一等待,避免提前 join 退化成串行;第三,每个下游任务要有超时和异常兜底;第四,区分核心依赖和非核心依赖,不能把订单主信息失败兜底成成功;第五,监控线程池 active、queue、拒绝次数和子任务 P95/P99,排查时结合 jstack 看线程是否卡在下游 IO、数据库连接或 Future.get。关联知识点
| 知识点 | 说明 |
|---|---|
| CompletableFuture 基础 | 基本 API 用法 |
| Callable、Future 与 FutureTask | JDK 7/8 异步结果模型对比 |
| 线程池总览 | 线程池是异步任务真正执行者 |
| 线程池生命周期与线上排查 | 线程堆积、Future.get、jstack 排查 |
| JDK 版本差异 | JDK 7 Future 到 JDK 8 CompletableFuture |
本章小结
CompletableFuture 的价值是表达异步依赖关系,不是让系统无限变快。它能把串行调用改成并行聚合,也能把依赖任务串成清晰链路;但如果线程池、超时、异常、下游容量和监控没设计好,异步只会把问题藏到后台线程里。真正会用 CompletableFuture,要同时懂任务编排、线程池容量、异常传播、超时取消和商业降级策略。
