欢迎光临
我们一直在努力

JAVA笔记之高级又好用的CompletableFuture

CompletableFuture 是 Java 8 引入的异步编程工具,实现了 Future 和 CompletionStage 接口。它支持链式组合、异常处理、多任务协调,是 Java 异步编程的核心类。

1. 基本概念

概念说明

Future

代表一个异步计算的结果,但只能 get() 阻塞等待

CompletionStage

定义了链式回调 API(thenApply 等)

CompletableFuture

两者都实现,既可手动完成,也可链式编排

2. 创建 CompletableFuture

2.1 已完成的 Future(同步值)

CompletableFuture<String> cf = CompletableFuture.completedFuture("hello");

2.2 异步执行(无返回值)

CompletableFuture<Void> cf = CompletableFuture.runAsync(() -> {
System.out.println("异步任务");
});

2.3 异步执行(有返回值)

CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> {
return fetchDataFromDB();
});

2.4 手动完成

CompletableFuture<String> cf = new CompletableFuture<>();
// 在别处,比如另一个线程完成
cf.complete("result"); // 正常完成
cf.completeExceptionally(new RuntimeException("error")); // 异常完成

2.5 指定线程池

默认使用 ForkJoinPool.commonPool()。生产环境建议自定义:

ExecutorService executor = Executors.newFixedThreadPool(10);
CompletableFuture.supplyAsync(() -> doWork(), executor);

//同样支持指定jdk 21+虚拟线程
Executor virtualExecutor = Executors.newVirtualThreadPerTaskExecutor();
CompletableFuture.supplyAsync(() -> doWork(), virtualExecutor);

3. 链式转换

3.1 thenApply — 转换结果

CompletableFuture<Integer> cf = CompletableFuture
.supplyAsync(() -> "hello")
.thenApply(s -> s.length()); // String -> Integer

3.2 thenAccept — 消费结果(无返回)

cf.thenAccept(result -> System.out.println("结果: " + result));

3.3 thenRun — 不关心结果,只执行动作

cf.thenRun(() -> System.out.println("任务完成"));

3.4 带 Async 后缀 — 在另一线程执行

cf.thenApplyAsync(s -> s.toUpperCase()); // 用 commonPool
cf.thenApplyAsync(s -> s.toUpperCase(), executor); // 指定线程池

4. 组合多个 Future

4.1 thenCompose — 扁平化嵌套(避免 Future<Future<T>>)

CompletableFuture<User> userFuture = getUserId()
.thenCompose(id -> getUserById(id)); // 返回 CompletableFuture<User>

等价于 flatMap,不要用 thenApply 嵌套 Future。

4.2 thenCombine — 合并两个独立 Future

CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<String> f2 = CompletableFuture.supplyAsync(() -> "World");
CompletableFuture<String> combined = f1.thenCombine(f2, (a, b) -> a + " " + b);
// "Hello World"

4.3 allOf — 等待全部完成

CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> "A");
CompletableFuture<String> f2 = CompletableFuture.supplyAsync(() -> "B");
CompletableFuture<String> f3 = CompletableFuture.supplyAsync(() -> "C");
CompletableFuture<Void> all = CompletableFuture.allOf(f1, f2, f3);

all.join(); // 阻塞直到全部完成

// 注意:allOf 返回 Void,需单独取各 Future 结果
String r1 = f1.join();

String r2 = f2.join();

4.4 anyOf — 任一完成即返回

CompletableFuture<Object> any = CompletableFuture.anyOf(f1, f2, f3);

Object first = any.join(); // 最先完成的那个结果

5. 异常处理

5.1 exceptionally — 捕获异常并恢复

CompletableFuture<String> cf = CompletableFuture
.supplyAsync(() -> {
if (Math.random() > 0.5) throw new RuntimeException("fail");
return "ok";
}).exceptionally(ex -> "默认值"); // 异常时返回 fallback

5.2 使用handle 

cf.handle((result, ex) -> {
if (ex != null) return "error: " + ex.getMessage();
return result;
});

5.3 whenComplete — 旁路观察(不改变结果)

cf.whenComplete((result, ex) -> {
if (ex != null){
log.error("失败", ex);
}else{
log.info("成功: {}", result);
}
});

5.4 completeOnTimeout / orTimeout(Java 9+)

cf.orTimeout(3, TimeUnit.SECONDS); // 超时抛 TimeoutException

cf.completeOnTimeout("default", 3, TimeUnit.SECONDS); // 超时返回默认值

6. 获取结果

// 阻塞等待
String result = cf.get(); // 可抛 InterruptedException, ExecutionException
String result = cf.get(5, TimeUnit.SECONDS); // 带超时

// 不抛受检异常(异常时抛 CompletionException)
String result = cf.join();
String result = cf.getNow("default"); // 未完成则返回默认值

7. 完整实战示例

示例 1:并行查询后聚合

public CompletableFuture<OrderDetail> getOrderDetail(String orderId) {
CompletableFuture<Order> orderFuture =CompletableFuture.supplyAsync(() -> orderService.getOrder(orderId));

CompletableFuture<List<Item>> itemsFuture =CompletableFuture.supplyAsync(() -> itemService.getItems(orderId));

CompletableFuture<User> userFuture =CompletableFuture.supplyAsync(() -> userService.getUser(orderId));

return orderFuture.thenCombine(itemsFuture, (order, items) -> new OrderDetail(order, items, null)).thenCombine(userFuture, (detail, user) -> {
detail.setUser(user);
return detail;
}).exceptionally(ex -> {
log.error("查询订单详情失败", ex);
return OrderDetail.empty();
});
}

示例 2:串行依赖链

public CompletableFuture<Report> generateReport(String userId) {
return CompletableFuture.supplyAsync(() -> userService.findById(userId))
.thenCompose(user -> dataService.fetchData(user.getId()))
.thenApply(data -> reportBuilder.build(data))
.whenComplete((report, ex) -> {
if (ex == null) cache.put(userId, report);
});
}

示例 3:批量并行 + 超时

List<CompletableFuture<String>> futures = ids.stream()
.map(id -> CompletableFuture
.supplyAsync(() -> callRemote(id))
.orTimeout(2, TimeUnit.SECONDS)
.exceptionally(ex -> "TIMEOUT"))
.toList();

CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();

List<String> results = futures.stream()
.map(CompletableFuture::join)
.toList();

8. 最佳实践

  • 生产环境指定线程池,避免占满 commonPool 影响 parallelStream 等。
  • 用 thenCompose 而不是 thenApply 处理返回 Future 的场景。
  • 始终处理异常,至少加 exceptionally 或 handle,避免异常被吞掉。
  • 避免在链里阻塞(如 get()、join()),会破坏异步意义;只在最终边界阻塞。
  • allOf 后要自己取各 Future 结果,它只返回 Void。
  • 注意默认执行线程:thenApply 可能在调用线程或完成线程执行;需要隔离时用 thenApplyAsync。
  • Java 9+ 可用 orTimeout、completeOnTimeout 做超时控制。

 

赞(0)
未经允许不得转载:171主机测评 » JAVA笔记之高级又好用的CompletableFuture
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址