如何用 CompletableFuture 编排 Java 异步任务?
学完这篇,你能用 CompletableFuture 把「查订单、查用户、查库存」这类互相独立的调用并行跑起来,并正确处理串联、合并、异常和超时。下面按步骤来。
第一步:确认 JDK 版本,准备好线程池
这一步要拿到一个可运行的环境和一个专用线程池,因为 CompletableFuture 的默认线程池不适合跑业务。在终端执行 java -version,确认版本号是 1.8.0 及以上(orTimeout、completeOnTimeout 需要 JDK 9+,delayedExecutor 需要 JDK 9+)。IntelliJ IDEA 里在 File → Project Structure → Project → SDK 选择 JDK 17 即可。
ExecutorService pool = Executors.newFixedThreadPool(16);
注意:不给 Executor 时,CompletableFuture 用的是
ForkJoinPool.commonPool(),它的线程数默认是 CPU 核数减 1。一旦你在里面做 HTTP 调用或数据库查询这类阻塞操作,公共池会被占满,整个 JVM 里其他用公共池的代码一起卡死。务必显式传入自己的线程池。
第二步:用 supplyAsync 把方法变成异步任务
这一步把同步方法包装成 CompletableFuture,得到的结果是一个「还没完成但将来会有值」的对象。
CompletableFuture<Order> orderF =
CompletableFuture.supplyAsync(() -> queryOrder("A1001"), pool);
CompletableFuture<User> userF =
CompletableFuture.supplyAsync(() -> queryUser("U2001"), pool);
supplyAsync 有返回值,runAsync 没有返回值(返回 CompletableFuture<Void>)。两行代码执行完,两个查询已经在并行跑了,主线程没有阻塞。
第三步:串联后续处理,分清 thenApply 和 thenCompose
这一步让上一个任务的结果流向下一个任务。两者区别很关键:
thenApply(Function):把结果映射成新值,适合做纯计算,比如.thenApply(Order::getAmount)。thenCompose(Function):当你的函数本身又返回一个 CompletableFuture 时用它,否则会得到嵌套的CompletableFuture<CompletableFuture<T>>。
CompletableFuture<String> chain = orderF
.thenCompose(order -> CompletableFuture.supplyAsync(
() -> queryLogistics(order.getId()), pool))
.thenApply(logistics -> "物流状态:" + logistics);
第四步:合并多个任务,用 thenCombine 和 allOf
这一步把多个并行结果汇总成一份数据。只有两个任务时用 thenCombine:
CompletableFuture<String> merged =
orderF.thenCombine(userF, (order, user) -> order + " | " + user);
System.out.println(merged.join());
三个以上任务用 allOf,它返回 CompletableFuture<Void>,本身不带结果,需要回头从各个 future 上取:
CompletableFuture.allOf(orderF, userF, stockF).join();
Order o = orderF.join();
User u = userF.join();
anyOf 用于「谁先返回用谁」,返回类型是 CompletableFuture<Object>,取用时需要自己强转。
注意:
join()抛出的是非受检的CompletionException;get()抛的是受检的InterruptedException和ExecutionException。在 Lambda 里做合并逻辑时用join()更省事。
第五步:处理异常,别让异常静默丢失
这一步保证某个子任务失败时你有兜底值。三种写法按场景选:
CompletableFuture<String> cf = CompletableFuture
.supplyAsync(() -> { throw new RuntimeException("库存服务超时"); }, pool)
.exceptionally(ex -> "库存未知"); // 只处理异常,返回兜底值
exceptionally(fn):只在异常时触发,返回兜底值。handle((res, ex) -> ...):正常和异常都会走,能改变返回值。whenComplete((res, ex) -> ...):只能做记录(打日志),改不了返回值。
注意:如果既没写
exceptionally也没写handle,异常会被吞掉,只有在调用join()/get()时才会以CompletionException形式抛出。上线前建议统一在末尾挂一个whenComplete打日志。
第六步:加超时,避免请求永远挂着
这一步给任务设置时间上限。JDK 9+ 直接用:
CompletableFuture<String> withTimeout =
orderF.thenApply(Object::toString)
.orTimeout(2, TimeUnit.SECONDS)
.exceptionally(ex -> "查询超时");
JDK 8 没有 orTimeout,可以用 ScheduledExecutorService 在超时后手动调 completeExceptionally,或者升级到 JDK 11/17(目前主流的 LTS 版本)。
最后记得 pool.shutdown(),否则主线程结束后 JVM 不会退出。
小结
- 必须传入自定义线程池,别用
ForkJoinPool.commonPool()跑阻塞任务。 thenApply做值转换,thenCompose做扁平化,返回 future 的函数一定用后者。- 多任务汇总用
allOf(...).join()后再逐个join()取结果。 - 异常用
exceptionally/handle/whenComplete兜底,否则会被静默吞掉。 - 超时用 JDK 9+ 的
orTimeout,JDK 8 需自行调度补齐。 - 线程池用完调用
shutdown()。
原文链接:https://www.gj0.com/thread-1664.html
转载请注明出处并保留本声明;内容仅代表作者观点,与本站立场无关。若本文涉嫌侵权,请联系本站处理。