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 就能从“会用”变成“敢在线上用”。

标签
CompletableFuture异步Java8并发