java
技术随笔
CompletableFuture 异步编程实战:并行查询、任务编排与异常处理
你还在用 Future.get() 傻等、还在写“串行查三个表再拼装”的慢接口吗?Java 8 引入的 CompletableFuture 是异步编程的“瑞士军刀”:既能并行执行多个任务,又能把异步结果像流水线一样串联、合并、编排异常处理。本文用最少的理论、最多的代码示例,带你把它用进日常开发。
一、Future 的三个痛点
future.get()会阻塞当前线程,异步变成了“假异步”,本质还是串行等待;- 无法表达任务之间的依赖关系(A 完成后 B 用 A 的结果继续算);
- 异常处理麻烦:只能 try-catch 包住 get(),任务链上的异常难以统一兜底。
CompletableFuture 解决了上述全部问题:它基于回调而非阻塞,天然支持任务编排,并提供丰富的异常处理 API。
二、创建异步任务
// 无返回值(相当于把 Runnable 丢进线程池)
CompletableFuture.runAsync(() -> log.info("写日志"));
// 有返回值(相当于 Callable)
CompletableFuture<User> f1 = CompletableFuture.supplyAsync(() -> userService.findById(1L));
// 强烈建议显式指定线程池(默认用 ForkJoinPool.commonPool,容易与业务线程池互相拖累)
Executor pool = Executors.newFixedThreadPool(8);
CompletableFuture<User> f2 = CompletableFuture.supplyAsync(
() -> userService.findById(2L), pool);
三、串行编排:thenApply / thenAccept / thenRun
CompletableFuture<String> f = CompletableFuture.supplyAsync(() -> "Order-100")
.thenApply(orderNo -> orderNo + ":paid") // 上一个结果作为入参,返回新结果
.thenApply(paid -> "[" + paid + "]"); // 继续处理
// thenAccept:只消费结果不返回(最终节点常用)
f.thenAccept(result -> log.info("最终结果: {}", result));
// thenRun:连结果都不需要,只要求“前一步完成后执行”
f.thenRun(() -> log.info("全部完成,做善后"));
// 异步变体 thenApplyAsync 等:换线程池执行回调,注意线程切换与上下文传递
记忆:Apply=结果传下去且返回新结果;Accept=吃下结果不吐出来;Run=吃完就完事。
四、并行组合:thenCombine / thenCompose
Executor pool = Executors.newFixedThreadPool(8);
// 典型场景:并行查用户 + 查订单,等两者都好后拼装响应(thenCombine 双向合并)
CompletableFuture<User> fu = CompletableFuture.supplyAsync(() -> userDao.findById(1L), pool);
CompletableFuture<List<Order>> fo = CompletableFuture.supplyAsync(() -> orderDao.listByUid(1L), pool);
CompletableFuture<UserDetailVO> result = fu.thenCombine(fo, (user, orders) -> {
UserDetailVO vo = new UserDetailVO();
vo.setUser(user);
vo.setOrders(orders);
return vo;
});
// thenCompose:前一个异步结果决定后一个异步任务(扁平化,避免 CompletableFuture 套娃)
CompletableFuture<User> f = CompletableFuture.supplyAsync(() -> "uid:1", pool)
.thenCompose(uidStr -> CompletableFuture.supplyAsync(() -> userDao.findById(1L), pool));
对比:thenApply 返回普通值,thenCompose 返回新的 CompletableFuture,后者用于“任务依赖另一个异步任务”的场景,防止出现 CompletableFuture<CompletableFuture<T>> 套娃。
五、批量等待:allOf / anyOf
// 场景:首页需要并行查询 N 个接口数据,全部完成后一次性返回
List<CompletableFuture<String>> futures = new ArrayList<>();
for (String biz : bizList) {
futures.add(CompletableFuture.supplyAsync(() -> remoteApi.call(biz), pool));
}
// allOf:等待全部完成(配合 join 汇总每个结果)
CompletableFuture<Void> all = CompletableFuture.allOf(
futures.toArray(new CompletableFuture[0]));
CompletableFuture<List<String>> resultFuture = all.thenApply(v ->
futures.stream().map(CompletableFuture::join).collect(Collectors.toList()));
// anyOf:任一完成即可(如多个数据源谁先返回用谁)
CompletableFuture<Object> any = CompletableFuture.anyOf(
CompletableFuture.supplyAsync(() -> query("mysql"), pool),
CompletableFuture.supplyAsync(() -> query("cache"), pool));
六、异常处理:exceptionally / whenComplete / handle
// exceptionally:出现异常时提供兜底值(类似 catch 返回默认)
CompletableFuture<Integer> f = CompletableFuture.supplyAsync(() -> 1 / 0, pool)
.exceptionally(ex -> { // ex 为 ExecutionException 的 cause
log.error("计算失败,使用默认值", ex);
return -1;
});
// whenComplete:无论成败都会执行,可拿到结果与异常(类似 finally,不吞异常)
CompletableFuture.supplyAsync(() -> "ok", pool)
.whenComplete((res, ex) -> {
if (ex != null) log.error("失败", ex);
else log.info("成功: {}", res);
});
// handle:既能处理成功结果也能处理异常,可返回新值(最灵活)
CompletableFuture<Integer> h = CompletableFuture.supplyAsync(() -> 1 / 0, pool)
.handle((res, ex) -> ex == null ? res : 0);
一个常见坑:join()/get() 会把异常抛给当前线程,若在异步回调里没处理又没被后续捕获,异常可能被吞掉。规范做法:任务链末端用 whenComplete 记录日志,或对每个分支显式 exceptionally 兜底。
七、实战:把慢接口从 900ms 优化到 300ms
// 改造前:串行调用三个远程服务
public DetailVO getDetailBefore(Long id) {
User u = userApi.get(id); // 300ms
List<Order> os = orderApi.list(id); // 300ms
List<Coupon> cs = couponApi.list(id); // 300ms
return assemble(u, os, cs); // 共约 900ms+
}
// 改造后:三个互不依赖的调用并行执行
public DetailVO getDetailAfter(Long id) {
CompletableFuture<User> fu = CompletableFuture.supplyAsync(() -> userApi.get(id), bizPool);
CompletableFuture<List<Order>> fo = CompletableFuture.supplyAsync(() -> orderApi.list(id), bizPool);
CompletableFuture<List<Coupon>> fc = CompletableFuture.supplyAsync(() -> couponApi.list(id), bizPool);
return CompletableFuture.allOf(fu, fo, fc) // 三个一起等,总耗时≈最慢的那个≈300ms
.thenApply(v -> assemble(fu.join(), fo.join(), fc.join()))
.exceptionally(ex -> assembleFallback(id)) // 兜底:任一出错返回降级数据
.join(); // 只在这里阻塞一次拿最终结果
}
八、开发中的注意点
- 务必自定义线程池并合理设置拒绝策略:默认 commonPool 与其他框架共用,核心接口异步化容易互相拖垮;
- 别在里面搞“异步套异步”:不要在 supplyAsync 内部再 supplyAsync 而不 join,容易把任务丢到别的池子造成失控;
- 上下文传递问题:异步线程拿不到父线程的 ThreadLocal(如登录用户、traceId),需要手动传递或使用 TransmittableThreadLocal;
- 事务与异步的组合:@Transactional 绑定当前线程,事务方法内用 CompletableFuture 并行写库要特别注意(子线程各自开事务,不共享主事务);
- 超时控制:用
get(timeout, TimeUnit)或orTimeout()给外部调用加超时,避免线程池被慢调用占满。
一句话收尾:串行接口用 thenApply 串联,互不依赖用 allOf/thenCombine 并行,最后用 exceptionally 兜底、用 join 收口——CompletableFuture 就能从“会用”变成“敢在线上用”。