Java异步编程实战:CompletableFuture任务编排与线程池优化
1. 从业务痛点认识CompletableFuture1.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 。听起来很简单但很多人在实际项目中并不怎么用它觉得没返回值的东西用处不大。我一开始也这么想后来发现它适合做“触发型”任务比如清理临时文件、发送通知、预热缓存。看个最基础的用法CompletableFutureVoid 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() ); CompletableFutureVoid future CompletableFuture.runAsync(() - { // 业务任务 }, pool);上面用了Guava的ThreadFactoryBuilder来给线程起名字这支操作小但极其有效。线上排查问题时线程dump一看名字就知道是哪个池子在干活不然全是pool-1-thread-1定位问题会想摔键盘。2.2 supplyAsync与thenAccept有返回值的第一种姿势如果任务需要返回结果就要用supplyAsync。它接收一个Supplier 返回CompletableFuture 。我来写一个实际场景假设现在要查询一个用户的积分余额这个操作耗时大概300毫秒我不想阻塞主流程。CompletableFutureInteger 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则只消费结果不返回新值。CompletableFutureString 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。换句话说它把两个异步操作串成一条链但不会阻塞等待第一个任务完成而是等它完成以后自动把结果传下去。CompletableFutureLong userIdFuture CompletableFuture.supplyAsync(() - { return 10001L; }, pool); CompletableFutureListOrder orderFuture userIdFuture.thenCompose(userId - CompletableFuture.supplyAsync(() - orderService.queryByUserId(userId), pool) );注意thenCompose返回的CompletableFuture是扁平的不会出现CompletableFutureCompletableFutureList 这种嵌套结构。这一点和thenApply有本质区别如果我用thenApplyFunction里返回的也是一个CompletableFuture那么得到的类型就是嵌套的后面还要手动flatMap似的拆开非常别扭。所以看到“返回值是一个异步任务”的时候优先考虑thenCompose。3.2 thenCombine并行合并两个独立任务更多的时候两个任务之间没有依赖关系比如同时查商品基本信息和库存数量最后组装成商品详情。这类场景就可以用thenCombine。它的作用是把两个CompletableFuture的结果同时交给一个BiFunction完成合并。CompletableFutureProductInfo infoFuture CompletableFuture .supplyAsync(() - productService.getInfo(1001L), pool); CompletableFutureInteger stockFuture CompletableFuture .supplyAsync(() - productService.getStock(1001L), pool); CompletableFutureProductDetail 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就是用来处理两个任务之间的竞争它会将先完成的任务结果传给后续函数。示例CompletableFutureString backupFuture CompletableFuture .supplyAsync(() - queryFromBackup(1001L), pool); CompletableFutureString primaryFuture CompletableFuture .supplyAsync(() - queryFromPrimary(1001L), pool); CompletableFutureString fastestFuture primaryFuture.applyToEither(backupFuture, result - result);如果存在多个任务可以用anyOf。它接收一个CompletableFuture数组返回一个CompletableFuture从架构角度来说这种竞争模式很值得在缓存系统中借鉴。比如Redis还没打好缓存并发请求同时打到数据库你可以让其中一个请求去查库回填其他请求等待同一个Future避免缓存击穿。这部分在架构设计里叫“请求合并”CompletableFuture配合ConcurrentHashMap也能实现简版。4. 几个容易踩坑的编排细节4.1 allOf等所有任务完成注意返回值类型很多业务场景需要等多个任务全部完成再聚合结果。allOf就是干这个的。它接收一个CompletableFuture数组返回CompletableFuture 。为什么返回值是Void因为它不负责帮你合并每个任务的结果它的作用只是提供一个“所有任务都完成”的信号。你要自己去拿每个Future的结果。组合多个结果时我习惯先建一个数组然后循环提交任务最后配合join把结果收集起来ListCompletableFutureItem futures new ArrayList(); for (Long itemId : itemIds) { CompletableFutureItem future CompletableFuture.supplyAsync(() - itemClient.getItem(itemId), pool); futures.add(future); } CompletableFutureVoid allDone CompletableFuture.allOf( futures.toArray(new CompletableFuture[0]) ); allDone.join(); ListItem 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它接收上一个阶段的异常返回一个默认值或做降级处理。比如查积分失败返回0CompletableFutureInteger 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可以做统一处理CompletableFutureString 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. 第一步并行获取基础数据 CompletableFutureProductInfo infoFuture CompletableFuture .supplyAsync(() - productClient.getInfo(productId), detailPool); // 库存不依赖其他数据也可以并行 CompletableFutureStock stockFuture CompletableFuture .supplyAsync(() - productClient.getStock(productId), detailPool); // 收藏数据可选如果失败不拖累主流程 CompletableFutureBoolean favFuture CompletableFuture .supplyAsync(() - userClient.hasFavorite(userId, productId), detailPool) .exceptionally(ex - { log.warn(query favorite fail, productId{}, productId, ex); return false; }); // 2. 从基本信息中拿到店铺ID继续串行获取推荐列表 CompletableFutureListProductSku 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可用于直接查询评分也可以与推荐并行走 CompletableFutureScore scoreFuture infoFuture.thenApply(info - scoreClient.getScore(info.getProductId()) ).exceptionally(ex - { log.warn(query score fail, ex); return Score.empty(); }); // 4. 聚合 CompletableFutureVoid 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(商品详情查询失败); }还有一种更优雅的做法是给每个子任务都设置超时CompletableFutureStock 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-catchtry { 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硬凑。最后分享一个小技巧在服务启动时把自定义线程池的核心参数打印到日志里比如线程池名称、核心线程数、最大线程数、队列容量。下次线上出问题打开日志第一眼就能确认有没有用错池子。这个习惯救过我一次可能也会救你一次。