CompletableFuture
CompletableFuture 是 Java 8 提供的异步编排工具。它常用于并行调用多个接口、异步执行耗时任务、组合多个任务结果。相比手动创建线程和等待,CompletableFuture 更适合表达异步任务之间的依赖关系。
本页先讲基础 API。如果要深入理解线程池选择、非 Async/Async 回调线程、异常传播、allOf 失败策略、超时是否真正取消底层任务、接口聚合线上排查,继续看:CompletableFuture 全过程原理与商业异步编排。
为什么需要 CompletableFuture
假设一个订单详情页需要查询用户、订单、优惠券、物流信息。如果串行调用,每个接口耗时 200ms,总耗时可能接近 800ms。
flowchart TD
A["查询用户 200ms"] --> B["查询订单 200ms"]
B --> C["查询优惠券 200ms"]
C --> D["查询物流 200ms"]
D --> E["总耗时约 800ms"]如果这些查询互不依赖,可以并行执行。
flowchart TD
A["发起请求"] --> B["查询用户"]
A --> C["查询订单"]
A --> D["查询优惠券"]
A --> E["查询物流"]
B --> F["汇总结果"]
C --> F
D --> F
E --> FCompletableFuture 执行原理
CompletableFuture 可以理解为“一个将来会完成的结果”。它内部保存任务状态、结果、异常和后续依赖动作。
flowchart TD
A["提交异步任务"] --> B["线程池执行任务"]
B --> C["任务成功或失败"]
C --> D["完成 Future 状态"]
D --> E["触发后续 thenApply 或 thenAccept"]
E --> F["继续向后传播结果"]它不是创建了一个神秘的新执行模型,真正干活的仍然是线程池。CompletableFuture 负责把多个任务的完成关系组织起来。
基本使用
import java.util.concurrent.CompletableFuture;
public class CompletableFutureBasicDemo {
public static void main(String[] args) {
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
sleep(500);
return "user";
});
String result = future.join();
System.out.println(result);
}
private static void sleep(long millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}supplyAsync 有返回值,runAsync 没有返回值。
thenApply、thenAccept、thenRun
import java.util.concurrent.CompletableFuture;
public class CompletableFutureChainDemo {
public static void main(String[] args) {
CompletableFuture.supplyAsync(() -> "Tom")
.thenApply(name -> "Hello " + name)
.thenAccept(System.out::println)
.thenRun(() -> System.out.println("处理完成"))
.join();
}
}| 方法 | 作用 |
|---|---|
thenApply | 接收上一步结果并返回新结果 |
thenAccept | 接收上一步结果但不返回 |
thenRun | 不接收结果,也不返回 |
thenCompose 和 thenCombine
thenCompose 用于有依赖的异步任务,thenCombine 用于合并两个独立任务结果。
import java.util.concurrent.CompletableFuture;
public class ComposeCombineDemo {
public static void main(String[] args) {
CompletableFuture<String> detailFuture = queryUserId()
.thenCompose(ComposeCombineDemo::queryUserDetail);
CompletableFuture<String> mergedFuture = queryOrder()
.thenCombine(queryCoupon(), (order, coupon) -> order + " + " + coupon);
System.out.println(detailFuture.join());
System.out.println(mergedFuture.join());
}
private static CompletableFuture<Long> queryUserId() {
return CompletableFuture.supplyAsync(() -> 1001L);
}
private static CompletableFuture<String> queryUserDetail(Long userId) {
return CompletableFuture.supplyAsync(() -> "user-" + userId);
}
private static CompletableFuture<String> queryOrder() {
return CompletableFuture.supplyAsync(() -> "order");
}
private static CompletableFuture<String> queryCoupon() {
return CompletableFuture.supplyAsync(() -> "coupon");
}
}allOf 汇总多个任务
import java.util.concurrent.CompletableFuture;
public class AllOfDemo {
public static void main(String[] args) {
CompletableFuture<String> userFuture = CompletableFuture.supplyAsync(() -> "user");
CompletableFuture<String> orderFuture = CompletableFuture.supplyAsync(() -> "order");
CompletableFuture<String> couponFuture = CompletableFuture.supplyAsync(() -> "coupon");
CompletableFuture<Void> all = CompletableFuture.allOf(userFuture, orderFuture, couponFuture);
all.join();
String result = userFuture.join() + ", " + orderFuture.join() + ", " + couponFuture.join();
System.out.println(result);
}
}allOf 本身不直接返回每个任务结果,只表示所有任务都完成。结果需要从各个 future 中取。
异常处理
import java.util.concurrent.CompletableFuture;
public class FutureExceptionDemo {
public static void main(String[] args) {
String result = CompletableFuture.supplyAsync(() -> {
if (true) {
throw new RuntimeException("查询失败");
}
return "success";
}).exceptionally(ex -> {
System.out.println("异常:" + ex.getMessage());
return "default";
}).join();
System.out.println(result);
}
}常见方法:
| 方法 | 作用 |
|---|---|
exceptionally | 异常时返回兜底值 |
handle | 成功和异常都处理,并返回结果 |
whenComplete | 成功和异常都能观察,但通常不改变结果 |
线程池为什么必须显式指定
不指定线程池时,异步任务默认使用公共线程池。业务系统里不建议重要任务都挤在公共线程池里。
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class FutureExecutorDemo {
private static final ExecutorService IO_POOL = Executors.newFixedThreadPool(8);
public static void main(String[] args) {
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
// 模拟 IO 查询
return "result";
}, IO_POOL);
System.out.println(future.join());
IO_POOL.shutdown();
}
}生产项目不要直接用 Executors.newFixedThreadPool 后就不管,还需要考虑线程命名、队列长度、拒绝策略和监控。这里为了演示保持简单。
超时控制
Java 9 之后可以使用 orTimeout 和 completeOnTimeout。
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
public class FutureTimeoutDemo {
public static void main(String[] args) {
String result = CompletableFuture.supplyAsync(() -> {
sleep(3000);
return "slow result";
}).completeOnTimeout("timeout default", 1, TimeUnit.SECONDS).join();
System.out.println(result);
}
private static void sleep(long millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}远程调用必须有超时,否则异步任务也会堆积。
商业 Demo:订单详情并行聚合
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class OrderDetailAggregateDemo {
private static final ExecutorService POOL = Executors.newFixedThreadPool(8);
public static void main(String[] args) {
CompletableFuture<String> user = CompletableFuture.supplyAsync(() -> "用户信息", POOL);
CompletableFuture<String> order = CompletableFuture.supplyAsync(() -> "订单信息", POOL);
CompletableFuture<String> logistics = CompletableFuture.supplyAsync(() -> "物流信息", POOL);
CompletableFuture.allOf(user, order, logistics).join();
String detail = user.join() + "," + order.join() + "," + logistics.join();
System.out.println(detail);
POOL.shutdown();
}
}真实项目里要加超时、异常兜底、线程池隔离、日志追踪。
常见风险
| 问题 | 后果 | 建议 |
|---|---|---|
| 不指定线程池 | 公共线程池被阻塞 | 业务线程池隔离 |
| 异步任务里做长时间阻塞 | 线程池耗尽 | 设置超时和容量 |
| 忘记处理异常 | join 时抛包装异常 | 使用 exceptionally 或 handle |
| 任务依赖关系写错 | 结果提前读取或死等 | 区分 thenCompose 和 thenCombine |
| 无限并行调用下游 | 打爆外部接口 | 限流和线程池隔离 |
本章小结
CompletableFuture 用来表达异步任务和任务依赖。它适合接口聚合、并行查询、异步处理,但不是性能魔法。生产使用时必须显式线程池、处理异常、设置超时,并控制对下游的并发压力。
