Java异步编程实战:CompletableFuture任务编排与线程池优化
2026/9/24 22:07:30 网站建设 项目流程

1. 从业务痛点认识CompletableFuture

1.1 为什么疯狂回调还不够用

做Java后端的朋友,多半都写过这样的代码:一个接口需要查用户信息、查订单列表、查优惠券,三个数据源各不相同。以前我习惯用线程池加Future,把三个任务丢进去,然后一个个get(),看起来逻辑清晰,但你会发现get()是阻塞的,如果串行执行或者等待时间过长,接口RT直接飙上去。后来用了回调,用CountDownLatch或者Guava的ListenableFuture,代码开始变得碎,尤其是多个任务之间有依赖关系时,回调套回调,套到最后自己都看不懂。

真正让我决定彻底转向CompletableFuture的,是一次活动页接口的重构。那个接口要聚合十几个下游数据,有的数据之间相互独立,有的数据必须等另一个返回后才能继续算。用回调写了一个晚上,越写越乱,最后用CompletableFuture半天就理清了。它把“异步执行”“结果转换”“异常处理”“任务编排”这四件事统一收进一套API里,代码读起来像同步流程一样顺畅,这才是它最值钱的地方。

如果你去翻Java 8的官方文档,会发现CompletableFuture同时实现了Future和CompletionStage两个接口。Future大家熟,但CompletionStage才是精髓。它定义了一系列方法,用来描述“当一个阶段完成之后,下一个阶段做什么”。这套设计思路,本质上就是函数式编程里的管道思想。用生活类比来说,普通Future像你去餐厅点餐,拿了个号,然后站在柜台前死等,直到菜端到你手上;CompletableFuture则像你拿了号以后可以随便逛,餐厅做好菜会通过手机通知你,你还可以顺便告诉服务员“吃完这盘菜帮我上甜点”,整个流程都由餐厅调度,你只需要描述规则。

1.2 适合谁来学、能解决什么

这篇内容适合三类人:第一类是刚接触Java并发、被各种锁和线程池绕晕的新人,你可以把CompletableFuture当成一个更友好的并发工具箱,先学会用,再慢慢理解背后的调度机制;第二类是天天写接口聚合、数据编排的后端开发,这一类人应该是受益最大的,很多令人头秃的串行等待,换成CompletableFuture以后,性能提升立竿见影;第三类是准备面试的同学,Java面试题里CompletableFuture出现频率越来越高,光会背八股文没用,真能写出优雅编排代码,才是面试官想看到的。

它能解决的问题,核心就是三个字:编排。比如A任务和B任务可以同时跑,C任务必须等A和B都结束才能开始,D任务只要A、B、C中任意一个成功就可以继续。这些复杂的依赖关系,用CompletableFuture的thenCombine、applyToEither、allOf、anyOf等方法,几行代码就能描述清楚。再配合自定义线程池,你可以把不同任务分配到不同池子,避免互相干扰。这篇文章我会从runAsync入门,一路讲到任务编排进阶,每个阶段都会给出可运行的代码案例,保证你照着敲就能跑出效果。

2. runAsync与supplyAsync:异步任务的起点

2.1 runAsync没有返回值,但别小看它

CompletableFuture入门第一课,通常是runAsync。它接收一个Runnable对象,执行一个没有返回值的任务,返回的CompletableFuture 。听起来很简单,但很多人在实际项目中并不怎么用它,觉得没返回值的东西用处不大。我一开始也这么想,后来发现它适合做“触发型”任务,比如清理临时文件、发送通知、预热缓存。

看个最基础的用法:

CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { System.out.println("任务执行线程: " + Thread.currentThread().getName()); }); future.join(); System.out.println("主线程结束");

默认情况下,runAsync会使用ForkJoinPool.commonPool()来执行任务,也就是公共的ForkJoin池。这里就有一个最常见的坑:ForkJoinPool公共池的线程数默认是CPU核心数减1。如果所有异步任务都往这个池子里丢,一旦某个任务发生了阻塞,比如调了第三方接口、查了数据库,公共池的线程就会被占住,其他依赖公共池的任务全都会排队。所以在生产环境,我强烈不建议直接用默认池,至少要指定一个自定义线程池。

ExecutorService pool = new ThreadPoolExecutor( 4, 8, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(500), new ThreadFactoryBuilder().setNameFormat("biz-pool-%d").get(), new ThreadPoolExecutor.CallerRunsPolicy() ); CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { // 业务任务 }, pool);

上面用了Guava的ThreadFactoryBuilder来给线程起名字,这支操作小但极其有效。线上排查问题时,线程dump一看名字就知道是哪个池子在干活,不然全是pool-1-thread-1,定位问题会想摔键盘。

2.2 supplyAsync与thenAccept:有返回值的第一种姿势

如果任务需要返回结果,就要用supplyAsync。它接收一个Supplier ,返回CompletableFuture 。我来写一个实际场景:假设现在要查询一个用户的积分余额,这个操作耗时大概300毫秒,我不想阻塞主流程。

CompletableFuture<Integer> scoreFuture = CompletableFuture.supplyAsync(() -> { // 模拟耗时操作 try { Thread.sleep(300); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } return 9527; }, pool); // 干点别的... System.out.println("主线程继续执行"); // 拿到结果 Integer score = scoreFuture.get(2, TimeUnit.SECONDS); System.out.println("用户积分: " + score);

get()方法会抛出ExecutionException和InterruptedException,而且会阻塞当前线程。与其调用get(),我更推荐join(),因为join()只抛CompletionException,不需要在方法签名里显式声明,代码更干净。当然如果调用了join()还希望控制等待时间,可以使用get(long, TimeUnit)或join配合orTimeout,这一点后面再聊。

有了结果以后,接下来就是对结果进行处理。常用的方式有两种:thenApply和thenAccept。thenApply用于将结果转换成另一个值,thenAccept则只消费结果,不返回新值。

CompletableFuture<String> resultFuture = CompletableFuture .supplyAsync(() -> 9527, pool) .thenApply(score -> "当前积分: " + score) .thenApply(msg -> msg + ",感谢参与活动"); resultFuture.thenAccept(System.out::println);

这里thenApply把上一步的返回值变成了字符串,然后又拼了一段文案。整个过程是异步的,每一个thenApply都会生成新的CompletableFuture,它们之间形成了一条任务链。这条链就是在CompletionStage上挂接的,你可以把它想象成流水线,每个工位只干一件事,做完以后传给下一个工位。

3. 异步任务编排:组合、合并与选择

3.1 thenCompose解决“依赖前一个任务的返回值”的场景

实际业务中,最典型的一类编排是串行依赖:第一个接口返回的值,是第二个接口的入参。比如先获取用户ID,再用这个ID去查订单。用Future写这种场景,要么嵌套回调用在future.get()之后再去发起第二次异步调用,阻塞等待,体验很差。

CompletableFuture里的thenCompose就是专门干这个的。它的参数是一个Function,入参是上一个阶段的结果,返回值类型是CompletionStage。换句话说,它把两个异步操作串成一条链,但不会阻塞等待第一个任务完成,而是等它完成以后自动把结果传下去。

CompletableFuture<Long> userIdFuture = CompletableFuture.supplyAsync(() -> { return 10001L; }, pool); CompletableFuture<List<Order>> orderFuture = userIdFuture.thenCompose(userId -> CompletableFuture.supplyAsync(() -> orderService.queryByUserId(userId), pool) );

注意,thenCompose返回的CompletableFuture是扁平的,不会出现CompletableFuture<CompletableFuture<List >>这种嵌套结构。这一点和thenApply有本质区别:如果我用thenApply,Function里返回的也是一个CompletableFuture,那么得到的类型就是嵌套的,后面还要手动flatMap似的拆开,非常别扭。所以看到“返回值是一个异步任务”的时候,优先考虑thenCompose。

3.2 thenCombine并行合并两个独立任务

更多的时候,两个任务之间没有依赖关系,比如同时查商品基本信息和库存数量,最后组装成商品详情。这类场景就可以用thenCombine。它的作用是把两个CompletableFuture的结果同时交给一个BiFunction,完成合并。

CompletableFuture<ProductInfo> infoFuture = CompletableFuture .supplyAsync(() -> productService.getInfo(1001L), pool); CompletableFuture<Integer> stockFuture = CompletableFuture .supplyAsync(() -> productService.getStock(1001L), pool); CompletableFuture<ProductDetail> detailFuture = infoFuture.thenCombine( stockFuture, (info, stock) -> { ProductDetail detail = new ProductDetail(); detail.setInfo(info); detail.setStock(stock); return detail; } );

两个任务可以同时执行,所以整体耗时近似于耗时更长的那个。使用thenCombine时,两个任务默认是并行执行的,因为每个supplyAsync都会立即提交到线程池。如果你传入了自定义线程池,且池子大小足够,那就能真正跑满并行度。

还有一点要提,thenCombine的变体很多,比如thenAcceptBoth,它不返回新值,只对两个结果做消费;还有runAfterBoth,两个任务都完成以后执行一个Runnable,不需要拿结果。选哪个,就看你要不要往下传值。

3.3 applyToEither与anyOf:谁先返回就用谁

再来聊一个很有意思的编排:多个竞争任务,取最先完成的那个。最典型的场景是“超时降级”或者“多路冗余请求”。比如为了降低接口延迟,可以同时向两个数据中心发同样的查询请求,谁先回来用谁的。

applyToEither就是用来处理两个任务之间的竞争,它会将先完成的任务结果传给后续函数。示例:

CompletableFuture<String> backupFuture = CompletableFuture .supplyAsync(() -> queryFromBackup(1001L), pool); CompletableFuture<String> primaryFuture = CompletableFuture .supplyAsync(() -> queryFromPrimary(1001L), pool); CompletableFuture<String> fastestFuture = primaryFuture.applyToEither(backupFuture, result -> result);

如果存在多个任务,可以用anyOf。它接收一个CompletableFuture数组,返回一个CompletableFuture

从架构角度来说,这种竞争模式很值得在缓存系统中借鉴。比如Redis还没打好缓存,并发请求同时打到数据库,你可以让其中一个请求去查库回填,其他请求等待同一个Future,避免缓存击穿。这部分在架构设计里叫“请求合并”,CompletableFuture配合ConcurrentHashMap也能实现简版。

4. 几个容易踩坑的编排细节

4.1 allOf等所有任务完成,注意返回值类型

很多业务场景需要等多个任务全部完成,再聚合结果。allOf就是干这个的。它接收一个CompletableFuture数组,返回CompletableFuture 。为什么返回值是Void?因为它不负责帮你合并每个任务的结果,它的作用只是提供一个“所有任务都完成”的信号。你要自己去拿每个Future的结果。

组合多个结果时,我习惯先建一个数组,然后循环提交任务,最后配合join把结果收集起来:

List<CompletableFuture<Item>> futures = new ArrayList<>(); for (Long itemId : itemIds) { CompletableFuture<Item> future = CompletableFuture.supplyAsync(() -> itemClient.getItem(itemId), pool); futures.add(future); } CompletableFuture<Void> allDone = CompletableFuture.allOf( futures.toArray(new CompletableFuture[0]) ); allDone.join(); List<Item> items = futures.stream() .map(CompletableFuture::join) .collect(Collectors.toList());

这里有个细节:allDone.join()只会在所有任务都结束时返回,但如果有某个任务异常,它并不会在join()时抛出那个异常,而是所有任务结束以后才返回。你需要单独在每个future.join()里捕获CompletionException。所以收集结果的循环里,最好加上try-catch或者exceptionally,否则某个任务失败会导致整个流程全部失败。

4.2 异常处理:exceptionally、handle和whenComplete的区别

异步任务的异常处理和同步代码不一样,无法靠try-catch包住整个流程,因为每步都是独立的。CompletableFuture提供了三个方法:exceptionally、handle和whenComplete。

exceptionally类似于catch,它接收上一个阶段的异常,返回一个默认值或做降级处理。比如查积分失败,返回0:

CompletableFuture<Integer> future = CompletableFuture .supplyAsync(() -> { if (true) throw new RuntimeException("查询失败"); return 100; }, pool) .exceptionally(ex -> { System.out.println("发生异常: " + ex.getMessage()); return 0; });

handle则更像是“无论如何我都会处理”,它的参数是BiFunction,第一个参数是上一个阶段的结果,第二个参数是异常。如果正常返回,异常就是null;如果异常抛出,结果就是null。所以handle可以做统一处理:

CompletableFuture<String> result = CompletableFuture .supplyAsync(() -> "数据") .handle((res, ex) -> { if (ex != null) { return "降级文案"; } return res + "_处理完成"; });

whenComplete只做通知,不改变结果。它和handle的核心区别是,whenComplete返回的CompletionStage会保留原始结果或异常,不会把whenComplete里的返回值作为新结果。

在实际项目中,我习惯在任务链的末端统一挂一个exceptionally来做兜底,防止异常在链上传递时被静默吞掉。同时,在关键节点用whenComplete打印日志,这样代码可观测性会好很多。

4.3 别在任务里用同步阻塞整个线程池

这里的坑和线程池直接相关。CompletableFuture本身很轻,但如果你在异步任务里调了同步阻塞的第三方HTTP接口,线程池中所有线程都会被占住。比如自定义线程池核心线程数8个,你提交了20个任务,其中有10个线程都卡在下游接口上,后续的任务就只能排队。这不是CompletableFuture的问题,而是任务设计和线程池配置不匹配的问题。

我的经验是:IO密集型任务,线程池大小可以设置为核心数乘2甚至更高,因为线程大部分时间在等待;CPU密集型任务,线程数就设置在核心数附近,避免线程过多导致频繁上下文切换。压测时还要关注队列容量和拒绝策略。如果使用了CallerRunsPolicy,当队列满了以后,新任务会在调用线程中直接执行,虽然不会丢任务,但会让调用线程的RT突然升高,这一点要和团队提前对齐预期。

5. 进阶实践:一个完整的任务编排案例

5.1 场景设计:商品详情页组装

把前面所有内容串起来,我们设计一个更接近生产的场景:商品详情页接口。需要获取以下数据:

  • 商品基本信息(必选)
  • 商品库存(必选)
  • 用户是否已收藏(需要登录,可选)
  • 同店铺推荐商品列表(可选,失败不影响主流程)
  • 商品评分(依赖商品ID,同时需要等基本信息返回后才有商品ID)

整个流程的依赖关系是:商品基本信息、库存、是否收藏可以并行;推荐商品列表依赖店铺ID,而店铺ID来自基本信息;评分同样依赖商品ID,不过它可以在拿到商品ID后立刻开始,而不必等库存和收藏。

用CompletableFuture编排起来,代码结构会非常清晰。我来写一个核心版本:

ExecutorService detailPool = buildCustomPool(16, "detail-pool"); // 1. 第一步并行获取基础数据 CompletableFuture<ProductInfo> infoFuture = CompletableFuture .supplyAsync(() -> productClient.getInfo(productId), detailPool); // 库存不依赖其他数据,也可以并行 CompletableFuture<Stock> stockFuture = CompletableFuture .supplyAsync(() -> productClient.getStock(productId), detailPool); // 收藏数据可选,如果失败不拖累主流程 CompletableFuture<Boolean> favFuture = CompletableFuture .supplyAsync(() -> userClient.hasFavorite(userId, productId), detailPool) .exceptionally(ex -> { log.warn("query favorite fail, productId={}", productId, ex); return false; }); // 2. 从基本信息中拿到店铺ID,继续串行获取推荐列表 CompletableFuture<List<ProductSku>> recommendFuture = infoFuture.thenCompose(info -> CompletableFuture.supplyAsync(() -> productClient.listRecommend(info.getShopId(), 10), detailPool) ).exceptionally(ex -> { log.warn("query recommend fail", ex); return Collections.emptyList(); }); // 3. 商品ID可用于直接查询评分,也可以与推荐并行走 CompletableFuture<Score> scoreFuture = infoFuture.thenApply(info -> scoreClient.getScore(info.getProductId()) ).exceptionally(ex -> { log.warn("query score fail", ex); return Score.empty(); }); // 4. 聚合 CompletableFuture<Void> all = CompletableFuture.allOf( infoFuture, stockFuture, favFuture, recommendFuture, scoreFuture ); all.join(); ProductDetail detail = new ProductDetail(); detail.setInfo(infoFuture.join()); detail.setStock(stockFuture.join()); detail.setFavored(favFuture.join()); detail.setRecommend(recommendFuture.join()); detail.setScore(scoreFuture.join());

这段代码的优点很明显:依赖关系清晰,每个数据源都是独立的变量声明;异常降级局部化,哪个失败就在哪个阶段处理,不会出现“一个失败全盘崩溃”的情况;可读性强,后续维护时一眼就能看到整个接口的数据流。

5.2 如何控制整体超时

上面的代码直接all.join(),如果某个下游接口特别慢,会拖垮整个接口。生产环境必须加超时控制。CompletableFuture本身的join没有超时参数,但get有,还可以用orTimeout。注意orTimeout是Java 9加入的,如果你的项目还是Java 8,就需要靠get(timeout, unit)自己实现。

try { all.get(2000, TimeUnit.MILLISECONDS); } catch (TimeoutException e) { // 超时后,未完成的任务结果会被丢弃 throw new BizException("商品详情查询超时"); } catch (ExecutionException e) { throw new BizException("商品详情查询失败"); }

还有一种更优雅的做法,是给每个子任务都设置超时:

CompletableFuture<Stock> stockFuture = CompletableFuture .supplyAsync(() -> productClient.getStock(productId), detailPool) .orTimeout(500, TimeUnit.MILLISECONDS) .exceptionally(ex -> Stock.empty());

这样单个数据源超过500毫秒就会返回空库存,主流程不会因为某个接口慢而整体超时。缺点也要说出来:orTimeout触发后,任务还在线程池里可能继续执行,线程资源并不会立刻释放,只是结果不会再被使用了。所以不能完全依赖它解决线程堆积问题,还是要给线程池设置合理的最大线程数和队列上限。

5.3 自定义线程池的工厂方法

我把构建线程池的代码抽成一个方法,方便复用。这里用到的几个参数,都是我在多个项目里压测后调出来的经验值,你可以根据自己的服务调整:

public static ExecutorService buildCustomPool(int coreSize, String namePrefix) { ThreadFactory factory = r -> { Thread t = new Thread(r, namePrefix + "-" + r.hashCode()); t.setDaemon(false); return t; }; return new ThreadPoolExecutor( coreSize, coreSize * 2, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(1000), factory, new ThreadPoolExecutor.CallerRunsPolicy() ); }

我习惯把核心线程和控制并发度紧密关联,而不是设置几百个线程。经验法则是:如果你的异步任务有80%的时间都是在等待IO,核心线程数可以设置为50到100;如果任务的CPU消耗很重,线程数设置为系统核心数加1到2就够了。线程不是越多越好,太多线程会导致切换开销,反之会浪费并发能力。

6. 常见问题与排查技巧实录

6.1 无返回值的任务怎么取异常

很多初学者用runAsync提交任务,任务内部抛异常,但外层调用future.join()却什么都没有。问题就在于,runAsync返回的CompletableFuture 同样会记录异常,只是默认没人看。只要在链路上没有显式处理,异常就被吞掉了。

排查方法很简单:在最终使用的future上调用exceptionally或者whenComplete打印日志。如果只是想临时看异常,可以在join()外直接try-catch:

try { future.join(); } catch (CompletionException e) { log.error("异步任务异常", e.getCause()); }

注意,CompletionException的cause才是真正的业务异常,不要把外层异常当成根本原因去查。

6.2 线程名“ForkJoinPool.commonPool-worker-1”意味着什么

线上日志里如果看到这个名字,说明你用CompletableFuture时没有传自定义线程池,默认使用的是ForkJoinPool.commonPool()。公共池的问题在于,它是JVM级别共享的,其他框架比如并行流(parallelStream)也在用它,一旦你的异步任务阻塞,其他使用公共池的地方也会跟着卡。

我的建议是:代码评审时,凡是出现CompletableFuture.supplyAsync(Supplier)没有第二个参数的,都要提出来,统一改成传入自定义线程池。这算是一个简单但收益很高的治理项。

6.3 内存泄漏或线程堆积的排查经验

有一次我们一个接口偶尔超时,线程dump一看,detail-pool的线程数全部处阻塞状态,队列里堆了几千个任务。当时的第一反应是下游服务变慢了,但后面深挖发现,是高并发时任务被大量提交到线程池,而每个任务内部又在等待另一个异步任务的结果,形成了死等。

具体场景是:任务A在等待任务B的结果,但任务B被提交到同一个线程池,而线程池线程全部被A占满,B排队进不来,于是A永远等不到B。这本质上是线程池饥饿问题。解决方式有三类:

  • 将不同依赖层次的任务放进不同的线程池,防止相互占用。
  • 控制主流程并发度,比如用信号量限制同时进入编排逻辑的数量。
  • 对等待加上超时,避免无限期堵塞。

排查的时候,线程dump配上jstack,重点看“waiting on”的信息,能很快定位到等待关系。纸上谈兵不如实际做一遍,建议你在本地用小的线程池复现一次饥饿问题,亲手观察线程状态,比背十遍八股文都有用。

7. 说点操作层面的真心话

CompletableFuture这套API不是看一遍文档就会的。我在写第一个版本时也踩过坑,比如在thenApply里写了耗时逻辑导致响应变慢,比如在thenCompose里没有复用同一个线程池,后期维护时才发现各个阶段线程不统一。这些问题的根源,不是API不友好,而是没有把“异步任务之间的依赖关系”当作架构设计的一部分来思考。

如果你刚开始学,我建议先跑通runAsync和supplyAsync,然后从thenApply、thenAccept、thenCompose、thenCombine这几个方法开始,逐个写一个小DEMO,拖动方法名参数看看IDE的提示,慢慢就会形成肌肉记忆。遇到异常处理时,多写几个exceptionally和handle的对比用例,把结果打印出来,你会比我当年更快理解它们之间的差异。

另外,不要迷信CompletableFuture能解决所有高并发问题。它的强项是编排,不是替代线程池、消息队列或者分布式事务框架。一个接口内多个数据源的异步聚合,用它非常合适;但如果你要做跨服务的数据一致性、重试补偿,还是老老实实引入消息中间件或者状态机,别用CompletableFuture硬凑。

最后分享一个小技巧:在服务启动时,把自定义线程池的核心参数打印到日志里,比如线程池名称、核心线程数、最大线程数、队列容量。下次线上出问题,打开日志第一眼就能确认有没有用错池子。这个习惯救过我一次,可能也会救你一次。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询