ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

CompletableFuture异步编排:解决串行RPC接口性能瓶颈

CompletableFuture异步编排:解决串行RPC接口性能瓶颈 线上有个商品详情页接口高峰期响应时间一直在1.8秒上下晃悠排查下来发现是十几个下游RPC串行调用一个等一个每个80毫秒加起来就奔着1秒去了。当时第一反应是多开线程并行调用但JDK自带的Future要拿结果就得阻塞get编排起来还特别别扭十几个任务的依赖关系用Future写得像一坨面条。后来换成CompletableFuture重构同样的逻辑压到300毫秒以内代码还清爽了不少。CompletableFuture和线程异步编排这套东西看起来API多、后缀杂但真拆开来看就四个动作创建、串联、并联、兜底把这四件事吃透日常业务里的异步场景基本都能覆盖。这篇文章适合正在被串行调用拖慢接口的初中级Java开发也适合想把异步代码写得更规范的资深同学。整篇不堆概念我会把每个API背后的调度逻辑、线程池怎么配、异常怎么兜、踩过的坑都摊开讲。如果你只想要一份能直接抄的模板第4章的实战案例可以整套拿走如果你想搞清楚为什么这么写前面几章的原理拆解值得慢慢看。1. 从Future的痛点到CompletableFuture的设计取舍1.1 Future那套为什么在业务编排里不好用先说清楚问题出在哪。Future是JDK5就有的东西提任务用ExecutorService.submit()拿结果用future.get()。它的设计目标是提交一个任务稍后来取结果本质上是一个异步结果的占位符。但业务里的异步从来不是提交一个任务这么简单往往是A的结果作为B的入参C和D并行跑完再合并这时候Future的三个短板就暴露得很明显。第一个短板是取结果只能阻塞。get()是阻塞方法get(timeout, unit)还是阻塞只是加了超时。你想在主线程把十个Future的结果收齐就得循环get()主线程被牢牢占住等于把并行退化成并发提交串行等待。真正想要的回调式处理Future给不了。第二个短板是没有编排能力。两个任务要合并结果Future只能靠你在外面手动get两次再算没有原生的combine语义。任务之间有依赖比如B依赖A你也只能先get A再提交B写起来全是在主线程里串糖葫芦异步的意义被吃掉大半。第三个短板是异常处理很别扭。Future的异常是在get()时以ExecutionException的形式抛出来的你得在每次get外面套try-catch还得层层解包。批量任务里某个失败其他结果的收集逻辑要写得非常小心稍不留神就是一个任务失败拖垮整批。1.2 CompletableFuture靠什么解决这些问题CompletableFuture在JDK8引入实现了两个接口Future和CompletionStage。前者保证它能当普通Future用后者才是重点——CompletionStage是一个阶段化编排的抽象它定义了当一个阶段完成后接下来要做什么这一整套组合语义。每个方法比如thenApply不是执行任务而是注册一个回调返回一个新的CompletableFuture代表新阶段像链条一样一节节接下去。它内部维护一个依赖栈当前阶段完成时会调用postComplete把注册在这个阶段上的回调挨个弹出来执行。因为触发和执行是分离的完成这件事既可能由工作线程触发也可能由主线程调用complete()手动触发。这个设计理解起来有个生活类比CompletableFuture像流水线上的一节节工位你在设计阶段就把这个工位干完东西传给下一个工位的规则定好齿轮转起来之后自动流转不需要人守在工位旁等着拿货。这套机制带来的直接好处有三个。回调式无阻塞后续动作靠回调触发发起方不用原地等待。原生组合能力thenCombine、thenCompose、allOf这些方法把并行的语义内建进来。异常可传递异常会沿着链条向下传播用exceptionally或handle在任意节点兜住像try-catch一样自然。这三点正好补上Future的三个短板所以现在业务里的异步编排基本都用它。2. 核心API按功能分类拆解而不是按字母背2.1 创建阶段supplyAsync和runAsync怎么选创建任务只有两个入口方法。supplyAsync(SupplierU)有返回值适合要算出一个结果的场景runAsync(Runnable)没有返回值适合只执行一个动作的场景。选择标准很朴素后续要拿这个结果做事的用supplyAsync纯执行副作用的用runAsync。拿商品价格去查优惠、拿用户ID去查标签这种都要返回值一律supplyAsync。这两个方法都有重载版本接收一个Executor参数。不传Executor时默认用的是ForkJoinPool.commonPool()这是个坑我后面第3章会专门讲为什么业务代码里几乎不该用它。创建方法的调用是立即返回的任务在线程池里排队执行不会卡住调用方。有个细节值得注意supplyAsync里的lambda如果抛异常异常不会在调用点抛出而是被封存进返回的CompletableFuture等你后续链到它时才浮出水面。这就是所谓的异常延迟暴露写的时候要意识到异常不会当场炸得在链尾处理。2.2 串联阶段thenApply、thenAccept、thenRun的语义差异这三个是使用频率最高的串联方法区别在于入参和返回值。thenApply(Function)接收上一阶段结果返回一个新结果用于数据转换比如把订单对象转成订单DTO。thenAccept(Consumer)接收结果但没有返回值用于消费比如把结果写入缓存或打日志。thenRun(Runnable)既不看上一步结果也没有入参纯在上一阶段完成后执行一个动作比如更新一个状态标记。这三者都有带Async后缀的版本比如thenApplyAsync。不带Async的版本会用触发上一阶段的线程或调用当前方法的线程执行回调。这句话很关键如果上一步是supplyAsync在工作线程完成的那么thenApply的回调也大概率由那个工作线程执行不额外消耗线程。而带Async的版本会把回调重新丢回线程池。什么时候要用Async当下游回调里有阻塞操作、耗时计算或者你想明确控制回调跑在哪个线程池时。如果只是轻量转换不带Async反而更省线程减少一次调度开销。2.3 依赖与聚合阶段thenCompose和thenCombine的分工这两个名字很像语义完全不同是新手最容易用错的。thenCompose(FunctionT, CompletionStageU)解决的是嵌套展开问题。假设你第一步拿到用户ID第二步要用这个ID去异步查用户详情而查详情本身返回的又是一个CompletableFuture直接用thenApply会得到CompletableFutureCompletableFutureUser这种双层套娃。thenCompose会把外层和内层拍平成一层和Stream的flatMap是一个思路。thenCombine(CompletionStageU, BiFunction)解决的是并行合并问题。两个互相不依赖的任务并行跑都完成后把两个结果合并。比如商品基础信息和库存状态同时查查完拼成一个完整对象。它和thenCompose的本质区别是Compose是串行依赖后一个任务依赖前一个的结果Combine是并行独立两个任务本身没有依赖只是最后要合在一起。搞混了会导致要么白白串行拖慢时间要么依赖没准备好就往上拼。2.4 多任务聚合allOf和anyOf的返回值陷阱涉及三个以上任务的聚合就用allOf或anyOf。allOf(CompletableFuture?...)返回一个CompletableFutureVoid语义是全部完成。这里有个必须记住的点它的返回值是Void不携带各个任务的结果。所以你拿到Void之后还得遍历原始的各个CompletableFuture去join()取结果。因为所有任务此刻都已完成join()不会阻塞所以遍历取结果是安全的。anyOf语义是任意一个完成返回CompletableFutureObject谁先完成就返回谁的结果。取到的结果类型是Object需要自己转。它适合多个数据源谁快用谁的场景比如同时请求两个推荐服务谁先返回用谁的结果。要注意anyOf完成后其他任务并不会被取消它们还在跑只是结果被忽略了。如果要精确控制资源得自己配合cancel。3. 线程池选型和依赖配置这步决定稳定性3.1 为什么说默认的ForkJoinPool commonPool是个隐患很多示例代码不传Executor用的就是ForkJoinPool.commonPool()。它在小demo里没问题但放到生产环境有两个致命问题。第一它的默认并行度是CPU核数减一。一台8核机器commonPool只有7个线程。你的业务如果同时有几百个异步任务全挤在这7个线程里排队延迟直接爆炸。第二它是整个JVM共享的。你项目里任何一个用parallelStream或者默认CompletableFuture的地方都在抢这7个线程。一个耗时的阻塞任务混进来会把commonPool的线程全占住导致所有依赖它的业务一起卡死。真实事故我见过一次某段代码用parallelStream处理一个大列表每个元素里又调了一次没传线程池的supplyAsync去查数据库。结果是列表处理把commonPool占满异步查库的任务全部排队等线程接口大面积超时。所以业务代码里创建CompletableFuture一定要显式传入自己的线程池这应该是团队里写进规范的一条铁律。3.2 自定义线程池的参数怎么算计算线程数有个经典公式核心线程数 CPU核数 × (1 等待时间 / 计算时间)。这里的等待时间指线程阻塞在IO上的平均时长计算时间指真正跑CPU指令的时长。假设你的任务是调RPC平均耗时100毫秒其中真正CPU计算只有5毫秒其余95毫秒都在等网络那么等待/计算约为198核机器上8 × 20 160。这个公式给的是理论参考值实际要结合压测调整因为线程太多上下文切换和内存占用也会上来。线程池的队列和拒绝策略也不能随便填。如果用无界队列LinkedBlockingQueue默认容量Integer.MAX_VALUE最大线程数就形同虚设任务会一直往队列里堆堆到内存溢出为止等于把问题从拒绝服务变成慢慢耗死。稳妥的做法是用有界队列比如容量500到1000加上明确的拒绝策略。拒绝策略推荐CallerRunsPolicy队列满时让调用线程自己执行起到天然的背压限流作用避免任务被无声丢弃。下面是一份可以直接改参数用的配置ThreadPoolExecutor bizPool new ThreadPoolExecutor( // 核心线程数按公式估算后取压测值这里以RPC密集场景举例 32, // 最大线程数核心的2倍左右突发时扩容 64, // 空闲线程存活时间 60L, TimeUnit.SECONDS, // 有界队列避免任务无限堆积 new LinkedBlockingQueue(500), // 线程命名方便排查问题强烈建议 new ThreadFactoryBuilder().setNameFormat(biz-async-%d).build(), // 队列满时由调用线程执行形成背压 new ThreadPoolExecutor.CallerRunsPolicy() );线程命名这一步特别值得强调。默认线程名是pool-1-thread-1这种出问题时打线程栈全是这种名字你根本分不清是哪个业务池。给每个业务线程池起个有意义的前缀比如order-async-goods-async-排查超时、线程泄漏时能省大量时间。3.3 线程池隔离与上下文透传一个原则不同业务、不同下游的异步任务尽量用不同的线程池。把所有异步任务塞进一个大池子看起来省事实际上一个下游抖动引发的堆积会连锁影响其他所有业务。按依赖分组订单查库存用一个池、查用户用一个池、查推荐用一个池每个池独立配置容量和超时这就是基本的隔离思路。另一个高频坑是上下文透传。主线程里设置的MDC日志追踪ID、ThreadLocal里的用户身份在supplyAsync切到工作线程后就丢了因为ThreadLocal是线程私有的。解决方案是手动在同一线程内先取值作为参数传进异步任务在任务内部再set回MDC。如果不想每个方法都手动做可以用TransmittableThreadLocal配合包装的线程池工厂在任务提交时把父线程的上下文快照带上执行时还原。日志追踪断了会让链路排查变成瞎子摸象这个投入很值。4. 实战把串行接口改造成并行编排4.1 需求拆解与依赖关系梳理拿一个典型的商品详情页举例。这个接口要返回商品基本信息、库存状态、当前用户对该商品的收藏状态、可用优惠券、推荐商品列表。改造前全是串行查商品 → 查库存 → 查收藏 → 查优惠券 → 查推荐五次RPC平均每次80毫秒加上一些本地处理接口耗时400毫秒起步高峰期放大到1.8秒。先做依赖分析商品基本信息是其他几项的前置因为库存、优惠券、推荐都要用到商品ID和类目ID而ID来自商品查询。所以结构是先查商品拿到基础信息后库存、收藏、优惠券、推荐这四项互相不依赖可以并行。这是个典型的一个串行头部 一段并行扇出 最终合并的编排模型。把依赖画清楚再写代码能避免把本可并行的任务误写成串行。4.2 编排代码实现与超时兜底按上面的结构代码大致长这样我把它拆成几步说明。第一步用supplyAsync查商品基础信息指定业务线程池。第二步用thenCompose衔接因为后续并行任务需要用到第一步的结果而并行段落本身要返回一个聚合后的CompletableFutureCompletableFutureItemDetailVO detailFuture CompletableFuture.supplyAsync(() - itemService.queryBase(itemId), bizPool) .thenCompose(base - { // 拿到base后四个任务并行发起互相不依赖 CompletableFutureStockVO stockF CompletableFuture.supplyAsync(() - stockService.query(itemId), bizPool) .exceptionally(e - StockVO.unknown()); // 单点降级 CompletableFutureBoolean favF CompletableFuture.supplyAsync(() - favService.isFav(userId, itemId), bizPool) .exceptionally(e - false); CompletableFutureListCouponVO couponF CompletableFuture.supplyAsync( () - couponService.query(userId, base.getCategoryId()), bizPool) .exceptionally(e - Collections.emptyList()); CompletableFutureListItemVO recommendF CompletableFuture.supplyAsync( () - recommendService.query(base.getCategoryId()), bizPool) .exceptionally(e - Collections.emptyList()); // allOf 聚合等四个都完成 CompletableFutureVoid all CompletableFuture.allOf( stockF, favF, couponF, recommendF); return all.thenApply(v - { ItemDetailVO vo new ItemDetailVO(); vo.setBase(base); vo.setStock(stockF.join()); vo.setFav(favF.join()); vo.setCoupons(couponF.join()); vo.setRecommends(recommendF.join()); return vo; }); }) // 整体兜底超时防止某个下游卡死拖垮接口 .orTimeout(500, TimeUnit.MILLISECONDS) .exceptionally(e - ItemDetailVO.defaultVO());这段代码有几个设计点值得展开。每个并行任务后面都挂了exceptionally做单点降级库存查失败就返回未知库存推荐查失败就返回空列表保证一个下游抖动不会让整个详情页打不开。这是局部降级优于全局失败的思路详情页这种场景少一块内容比整个页面报错体验好得多。allOf的返回值是Void所以我在thenApply里用各个任务的join()取结果。这里用join()而不是get()区别是get()会抛受检异常InterruptedException和ExecutionException逼着你写try-catch而join()抛的是非受检的CompletionException在链式代码里更顺手。因为此刻allOf已经保证了全部完成join()不会阻塞直接拿值。最后挂了orTimeout(500, MILLISECONDS)做整体超时保护。这个超时兜底非常关键即使前面的单点降级都失效、某个下游RPC挂死整体也能在500毫秒后返回兜底VO接口不会无限等待。没有这层保护一个卡死的下游就能把整个接口拖垮。4.3 异常语义与降级策略的配合这里要理清exceptionally在整个链里的作用位置。它只处理它前面那段链的异常并且只有当异常发生时才触发。上面每个并行任务单独挂exceptionally是为了让该任务的失败被就地消化返回一个中性值这样allOf层面看到的是全部成功不会因为一个任务失败而整体走异常分支。这是把异常控制在最小范围的实践。如果你希望某个关键任务失败时整个接口失败就不要给它挂exceptionally让异常往上冒在最外层用handle或exceptionally统一兜底。降级策略要根据业务重要性分级商品基础信息查不到详情页就没意义了应该整体失败或跳错误页推荐列表查不到返回空即可。把每个下游的失败代价想清楚再决定是在单点兜还是整体兜这比无脑全兜或者全不兜都更合理。5. 异常处理和超时控制兜底做不好前面全白搭5.1 handle、exceptionally、whenComplete的适用场景这三个方法经常被混用弄清它们的差异能少写很多冗余代码。handle(BiFunctionT, Throwable, U)无论成功失败都会执行入参是结果或异常返回值会作为新阶段的结果。它适合成功和失败都要处理且要产出新结果的场景比如成功返回数据、失败返回默认值两种路径都要走同一段转换逻辑。因为入参既有结果又有异常用的时候要先判断哪个是null。exceptionally(FunctionThrowable, T)只在失败时执行成功时它整个被跳过结果原样透传。它适合纯粹的异常兜底出错了给个默认值没出错别打扰。像上面详情页里每个单点降级用它就很合适语义干净。whenComplete(BiConsumerT, Throwable)无论成功失败都执行但没有返回值而且它不会改变链条上的结果——成功还是那个结果失败还是那个异常继续往下传。它适合做副作用比如记录日志、上报监控、释放资源做完之后不影响主流程。选型口诀要转换结果用handle只在失败时兜底用exceptionally只做副作用不影响流程用whenComplete。5.2 orTimeout和completeOnTimeout的区别超时控制有两个方法。orTimeout(long, TimeUnit)在超时后以TimeoutException异常结束这个CompletableFuture。completeOnTimeout(T value, long, TimeUnit)在超时后用给定的默认值正常完成。区别就是一个走异常路径、一个走成功路径。选哪个看你的下游处理习惯如果外层已经有统一的异常兜底逻辑用orTimeout让超时当作异常处理更一致如果希望超时直接给个值、不想走异常分支用completeOnTimeout更直接。有个容易踩的坑这两个方法的超时是作用于调用它的那个CompletableFuture及其上游链不能单独约束某一个上游任务。比如你在合并后的future上挂orTimeout它是从挂载那一刻开始计时到整个链完成。如果想给每个并行任务单独设超时得在各自的future上分别挂。还有一种情况超时触发后上游任务并不会被中断它还在后台跑只是结果被忽略了。如果上游是阻塞IO、会长期占线程光靠超时不够还要在任务内部配合可中断的调用或响应中断。5.3 异常被静默吞掉的几种典型情况最隐蔽的一类bug是异常被吞掉任务看起来成功实际啥也没干。常见场景有几种。第一种是链上某处挂了exceptionally返回了null后面接着用这个结果NPE被后续的catch又兜了一次最终表现为数据缺失而不是报错。第二种是allOf之后直接join()某个任务异常会在join时以CompletionException形式抛出如果这层没处理异常会冒到最外层线程日志里只有一条异步线程异常业务层毫不知情。第三种是任务提交后根本没消费返回的CompletableFuture异常永远没机会暴露任务默默失败。排查这类问题的思路是给每个异步任务补上日志和监控。在task里加try-catch打错误日志或者在链尾挂whenComplete记录状态。不要指望异常自己会冒出来异步世界的异常是你不主动接它就消失。生产环境里建议对关键异步任务加埋点统计成功率、耗时分布异常率一抬头立刻能发现。6. 常见问题排查与速查表6.1 线程池打满和任务堆积的排查路径异步任务出问题第一个要看的永远是线程池状态。排查路径是先看线程池的活跃线程数、队列长度、完成任务数。如果活跃线程长期等于最大线程数队列还在涨说明线程池被塞满了这时候要么是并发量突增、容量不够要么是某个任务卡住了拿住线程不放。前者扩容量后者找卡点。任务卡住的原因通常有几个下游RPC没有设超时一直在等任务里有死循环或长计算任务之间形成了循环依赖导致互相等。下游调用必须设超时这是异步场景的底线因为线程是共享资源一个任务无限期占用会连累所有排队的任务。用jstack导出线程栈找biz-async-前缀的线程看它们在哪个方法上停着能快速定位卡点。6.2 死锁和线程饥饿的预防CompletableFuture里最隐蔽的死锁是任务在自己的线程池里提交子任务并等待子任务完成。比如线程池只有4个线程你提交了4个任务每个任务内部又用同一个池提交了一个子任务并阻塞等它完成结果4个线程全在等子任务子任务却在队列里排不到线程死锁形成。这类问题叫线程池线程耗尽型死锁排查时线程栈会显示一堆线程阻塞在join()或get()上而队列里还有没执行的任务。预防办法有两个方向。一是避免在池内任务里用同一个池提交并等待子任务如果一定要有层级任务给子任务用不同的池或者控制好池的容量远大于可能的并发等待数。二是直接不用阻塞等待用thenCompose或thenCombine把依赖关系交给编排链让子任务完成后自动触发后续而不是靠线程阻塞去等。这个思路的本质是把线程等待换成回调驱动也是CompletableFuture相比Future更先进的地方。6.3 常见问题速查表现象可能原因排查方向处理建议接口偶发超时默认commonPool线程不够看线程栈是否有ForkJoinPool.commonPool显式传入业务线程池异步任务结果丢失未消费返回的future异常被吞检查任务是否被正确链式消费补whenComplete日志和埋点队列持续增长线程池容量不足或任务卡住监控活跃线程、队列长度扩容或定位卡点任务线程栈大量阻塞在join池内任务等待同池子任务jstack看等待链拆分线程池或用编排替代阻塞日志追踪ID断链ThreadLocal未跨线程透传看异步任务内MDC是否为空手动透传或用TTL包装超时后任务还在跑超时不会中断上游确认上游是否有可中断机制任务内响应中断或加内部超时7. 一些写在规范里的实践建议最后分享几条我踩过坑之后固化到团队规范里的东西都是些不起眼但很影响稳定性的细节。创建CompletableFuture永远显式传线程池一条都不能例外这条能规避掉一大半和commonPool相关的诡异问题。每个线程池必须有名字前缀从命名上就能看出归属故障排查时价值极大。所有下游调用都要有超时不管是RPC、数据库还是缓存异步环境里没有超时的调用就是定时炸弹。再补一个关于join()和get()选择的细节在编排链内部我用join()居多因为它抛非受检异常代码干净但在需要精确处理中断、或者处于阻塞敏感的线程里用get()更合适因为get()能响应InterruptedException。这个选择没有绝对对错看你这段代码对中断信号是否敏感。还有一点是关于测试。异步代码的测试不能只跑一遍看结果对不对要专门构造异常路径和超时路径验证降级逻辑真的生效。我见过太多降级代码从来没被真正执行过等到线上真出事才发现降级分支本身就是错的。用Mockito让某个下游抛异常或者延迟返回跑一遍看最终返回的是不是预期兜底值这个测试用例的成本很低收益却很高。我现在做异步改造的固定流程是先画依赖图理清串并行再按下游划分线程池然后写编排链、逐点挂降级和超时最后补日志埋点和异常路径测试。这套流程走下来异步代码的稳定性基本可控剩下的就是随业务量调整线程池参数了。
返回列表