CompletableFuture 异步编排实战:线程池、超时与异常处理
发布时间:2026/10/1 22:50:13来源:尧图网络
第一次在生产环境里用 CompletableFuture 把五个串行接口压成一次并行调用接口 P95 从 480ms 掉到 130ms那种感觉确实挺爽。但同样也是它让我在某个凌晨两点对着一个永远不返回的 join() 干瞪眼最后发现是自定义线程池的队列满了、拒绝策略还是 CallerRunsPolicy主线程被自己提交的任务给堵死了。CompletableFuture 这个工具就是这样上手半小时就能写出能跑的代码但要写出在生产环境里扛得住、查得动、改得动的异步编排中间隔着好几层坑。这篇东西我打算把它讲透从它到底补上了 Future 的哪些短板到 thenApply、thenCompose、thenCombine 这几个最容易混的 API 该怎么选再到线程池参数怎么算、超时怎么设、异常为什么老是凭空消失、线上线程池打满该怎么排查。文中的代码片段都是可以直接拿去改的骨架参数计算过程我会一步步写出来不看结论也能自己推。适合已经会写 Java、但对异步编排还停留在抄一段能用就行阶段的同学也适合那些代码已经上了生产、正在被莫名其妙的超时和线程堆积折磨的同学。1. 为什么值得把 CompletableFuture 当成主力异步编排工具1.1 从 Future 那点憋屈说起Java 5 就把 Future 塞进标准库了初衷很清楚把耗时操作丢给别的线程主线程该干嘛干嘛。但真正在业务代码里用过的人心里都有数Future 的 API 少得可怜——提交、取消、判断是否完成、get 拿结果就这四件事。它的定位其实是一个结果占位符压根就不是编排工具。这个定位差异直接导致了它在稍微复杂一点的场景下就力不从心。举个具体的例子。你在做一个商品详情页需要查商品基础信息、库存、价格、评价数量、推荐列表这五份数据来自五个不同的下游服务。用 Future 你只能 submit 五次然后五次 get。这个 get 是阻塞的谁先完成谁后完成你控制不了只能按固定的顺序一个个等。五个接口平均各 80ms最慢的 200ms总体耗时就是五者之和再叠上网络抖动串行的本质没变你只是把等待挪到了不同的 get 调用上而已。更要命的是依赖关系。假设推荐列表的请求参数里需要带上商品类目而类目要从基础信息里取这就变成了两段式调用先拿基础信息再拿推荐。用 Future 你只能等第一个 get 返回拿到类目再 submit 第二个任务再加到列表里等。代码写成什么样子可以想象一堆中间变量、一堆 try-catch、一堆手动维护的 List 可读性基本为零。异常处理也是老问题。Future.get() 抛出的是 ExecutionException真正的业务异常被包了一层你得 catch 住之后再调用 getCause() 才能拿到原始异常如果中间还有别的框架又包了一层那就得剥好几层。日志里打出来的栈信息经常是一堆 java.util.concurrent.ExecutionException翻半天看不到真正的报错点。还有一个很隐蔽的限制Future 没法手动完成。有些场景下结果不是由计算产生的而是由外部事件触发的比如一个回调、一条消息、一次信号通知。这种场景你用 Future 就只能开一个线程去阻塞等待白白占着一个线程不放线程池稍微小一点就直接被占满了。这就是典型的用阻塞的方式实现异步听起来荒谬但在 Future 时代确实是常规操作。1.2 CompletableFuture 补上的三块拼图CompletableFuture 从 Java 8 进入标准库它同时实现了 Future 和 CompletionStage 两个接口。这两个接口恰好对应了两层用途Future 那一层负责最终拿到结果CompletionStage 那一层负责把多个阶段串起来、拼起来、合起来。理解了这层双身份很多 API 的设计意图就通了。它补上的第一块拼图是回调式编程。thenApply、thenAccept、thenRun 这一族方法的作用是在当前阶段完成的那一瞬间自动触发后续动作完全不需要有线程在那里阻塞等待。这一点带来的直接好处是线程利用率。一个线程提交完任务就可以立刻回去处理别的请求等结果就绪时由完成该任务的线程顺手把回调执行掉或者异步交给另一个线程池执行。整个过程没有线程被空转占用。第二块拼图是组合能力也是它真正被称为编排工具的原因。thenCompose 用来串联两个有前后依赖的异步操作把 CompletableFutureCompletableFuture 这种嵌套结构扁平化thenCombine 用来把两个彼此独立的异步结果合并成一个allOf 和 anyOf 用来聚合一批任务。有了这几个方法前面那个先查基础信息再查推荐的两段式调用就能写成一个连贯的链式表达式。第三块拼图是手动完成complete() 和 completeExceptionally() 这两个方法让 CompletableFuture 成了一个可以被外部驱动的状态容器超时控制、事件驱动、回调桥接这些场景都靠它实现。1.3 什么场景适合用什么场景别硬上我不建议把 CompletableFuture 当成银弹到处套。异步带来的复杂度是实打实的线程上下文会丢、异常栈会变得又长又难读、调试的时候断点跳来跳去、线程池管理不当会直接拖垮整个应用。一个只有两个串行调用、总耗时 120ms 的接口改成并行之后变成 90ms省下来的 30ms 根本抵不上你为它引入的维护成本。我的判断标准大概是这样。适合用的场景有三个特征调用的下游数量多三个以上、彼此之间没有强依赖或者依赖关系是清晰的树状、单个调用的耗时有明显波动说明有 I/O 等待可以榨取。典型的例子就是详情页聚合、订单列表批量补数据、报表多维度统计、风控规则并行执行。这些场景下并行化带来的收益是成倍的代码复杂度换来的性能提升划算。不适合的场景也很明确。CPU 密集型任务本来就该用并行流或者专门的线程池套 CompletableFuture 只会让代码更绕。完全串行依赖的流程比如扣款成功才能发货异步化毫无意义反而把事务边界搞乱了。还有一种容易被忽略的情况是低频接口比如一天调用几百次的运营后台功能为了它去折腾异步编排收益基本可以忽略写清楚比写快更重要。2. 核心 API 全景拆解分清每个方法到底在干什么2.1 创建任务supplyAsync 和 runAsync 的差别不只是有没有返回值CompletableFuture 的入口方法主要就四个重载runAsync(Runnable)、runAsync(Runnable, Executor)、supplyAsync(Supplier)、supplyAsync(Supplier, Executor)。前两个没有返回值后两个有返回值这个差别谁都看得出来。真正值得说的是另一个差别不传 Executor 的时候任务会跑到 ForkJoinPool.commonPool() 里而 commonPool 的默认线程数是 CPU 核数减一。这意味着在一台四核的机器上你的所有没指定线程池的 CompletableFuture 任务都挤在三个线程里抢。一旦链路上有任何一个环节发生了阻塞比如一次没设超时的 HTTP 调用这三个线程很快就被占满然后所有依赖 commonPool 的代码一起卡住。这不是理论风险我在项目里见过不止一次某个同事写了个工具类里面用 supplyAsync 没传线程池量一上来整个应用的异步链路全堵住日志里看不到任何报错就是响应越来越慢。所以我给自己定的规矩很硬业务代码里的 supplyAsync 和 runAsyncExecutor 参数一律不能省。哪怕只是临时用一下也要显式传一个自己管理的线程池哪怕就是 Executors.newFixedThreadPool(4) 这种简单实现至少它和 commonPool 是隔离的出问题的时候边界清楚。// 不推荐默认跑在 commonPool线程数受限且与全应用共享 CompletableFutureString bad CompletableFuture.supplyAsync(() - queryRemote()); // 推荐显式指定业务隔离的线程池 Executor bizPool new ThreadPoolExecutor( 8, 16, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue(256), new NamedThreadFactory(detail-fetch-), new ThreadPoolExecutor.CallerRunsPolicy()); CompletableFutureString good CompletableFuture.supplyAsync(() - queryRemote(), bizPool);提示commonPool 并不是完全不能用。纯 CPU 计算、没有阻塞、任务粒度小且数量可控的场景用它其实是合适的。关键在于你要清楚自己在用哪个池子而不是忘了传。2.2 结果转换thenApply、thenAccept、thenRun 三个门派的边界这三兄弟是链式写法里出现频率最高的但它们的分工完全不同混用会导致编译不过或者行为不符合预期。判断方法很简单看输入和输出。thenApply(FunctionT,R)拿到上一步的结果处理之后返回一个新结果。它是转换会改变阶段的值。适合做数据加工比如把远程返回的 JSON 字符串转成对象、把金额单位从分转成元。thenAccept(Consumer )拿到上一步的结果消费掉但不返回值。它的返回类型是 CompletableFuture 。thenRun(Runnable)连上一步的结果都不需要只是在前面完成后触发一个动作同样返回 CompletableFuture 。这三者还有一个共同的后缀版本thenApplyAsync、thenAcceptAsync、thenRunAsync。带 Async 的方法会把后续动作交给线程池执行不带 Async 的方法则在完成前一个阶段的那个线程上直接执行。这一点非常关键也是最容易踩坑的地方。假设你的链路是这样supplyAsync(fetch, poolA).thenApply(this::parse)。如果 fetch 是在 poolA 的线程 X 上完成的那么 parse 也会在 X 上执行。如果 parse 是个耗时 50ms 的 JSON 解析那 X 就被占用了 50ms 不能干别的。更极端的情况是如果你的 parse 里有阻塞调用那就是在浪费线程池的宝贵线程。所以我的经验是轻量级、纯内存的转换用不带 Async 的版本省一次线程切换性能反而更好任何可能阻塞或者耗时不确定的操作一律用带 Async 的版本并且显式指定线程池。CompletableFutureOrder future CompletableFuture .supplyAsync(() - fetchOrderRaw(orderId), ioPool) // I/O 在 ioPool .thenApply(this::parseOrder) // 内存解析线程复用即可 .thenApplyAsync(this::enrichWithPromotion, cpuPool) // 涉及计算切到 cpuPool .thenApply(order - { log.info(order assembled: {}, order.getId()); return order; });2.3 thenCompose 与 thenCombine编排里最容易混的一对这两个名字长得像、签名也像但语义差得远几乎是面试和 code review 里的必考题。thenCompose(FunctionT, CompletionStage)处理的是依赖关系。当你的下一步操作需要用到上一步的结果作为参数而且这个下一步本身又是异步的就用它。它的作用是扁平化把 CompletableFutureCompletableFuture 压成 CompletableFuture。thenCombine(CompletionStageother, BiFunctionT,U,V)处理的是并列关系。当你有两个互相独立的异步任务需要把两个结果合起来做一件事就用它。两个任务会同时跑谁先谁后不重要。用一个具体的例子区分。查订单详情的时候从订单服务拿订单基本信息和从用户服务拿用户资料这两个动作互不依赖可以并行最后一起组装成页面数据这就是 thenCombine 的场景。而先拿订单基本信息从里面取出用户 ID再拿用户资料这就是 thenCompose 的场景因为第二步的参数依赖第一步的产出。如果该用 thenCompose 的地方误用了 thenApply你会得到 CompletableFutureCompletableFuture 这种嵌套类型写起来很别扭而且外层 future 完成的时候内层任务可能还没跑完很容易出现拿到的是个没完成的 Future这种诡异 bug。反过来如果该用 thenCombine 的地方用成了顺序的 thenCompose功能上不会错但两个本该并行的调用被强制串行了性能白白损失一半。// 场景一并行 合并用 thenCombine CompletableFutureOrder orderF CompletableFuture.supplyAsync(() - orderService.get(orderId), ioPool); CompletableFutureUser userF CompletableFuture.supplyAsync(() - userService.get(userId), ioPool); CompletableFutureDetailVO vo orderF.thenCombine(userF, this::assembleDetail); // 场景二依赖 串联用 thenCompose CompletableFutureUser userF2 CompletableFuture .supplyAsync(() - orderService.get(orderId), ioPool) .thenCompose(order - CompletableFuture.supplyAsync( () - userService.get(order.getUserId()), ioPool));2.4 allOf / anyOf多任务聚合的真实行为当并行任务从两个变成五个、十个thenCombine 就不好写了因为它只接受两个参数。这时候要用 allOf。CompletableFuture.allOf(f1, f2, f3, ...) 返回一个 CompletableFuture 它在所有传入的任务都完成后才完成。注意它返回的是 Void也就是说它只告诉你都完事了但不携带任何结果数据。这是初学者最容易懵的地方——很多人以为 allOf 会返回一个结果列表结果发现类型是 Void不知道怎么取值。正确的用法是先保存每个子任务的引用用 allOf 等待全部完成然后从各个引用上逐个 join 取结果。CompletableFutureBaseInfo base CompletableFuture.supplyAsync(() - baseService.get(sku), ioPool); CompletableFutureStock stock CompletableFuture.supplyAsync(() - stockService.get(sku), ioPool); CompletableFuturePrice price CompletableFuture.supplyAsync(() - priceService.get(sku), ioPool); CompletableFutureComments cmt CompletableFuture.supplyAsync(() - cmtService.get(sku), ioPool); CompletableFutureVoid all CompletableFuture.allOf(base, stock, price, cmt); CompletableFutureDetailVO result all.thenApply(v - { DetailVO vo new DetailVO(); vo.setBase(base.join()); vo.setStock(stock.join()); vo.setPrice(price.join()); vo.setComments(cmt.join()); return vo; });这里每个 join() 都不会真的阻塞因为 allOf 已经保证了所有任务都已完成join 只是把结果取出来。这个细节很重要很多人担心 join 会阻塞其实在 allOf 的回调里它是立即返回的。注意如果其中任何一个子任务抛出了异常allOf 返回的 future 会立刻以异常完成不会等剩下的任务。而剩下的任务其实还在跑只是没人管了。这个行为在设计降级策略的时候要考虑到。anyOf 的行为是任意一个完成就返回它返回的是 CompletableFuture还有一个聚合的替代方案值得提一下如果你用的是 Java 9 及以上可以试试 CompletableFuture 的延迟组合或者干脆用 Stream API 配合自定义的 collector 做批量聚合。不过在 Java 8 环境里allOf 加 join 这个组合依然是最稳、最通用的写法。3. 线程池选型与参数计算实操3.1 默认的 commonPool 为什么不能当生产主力前面简单提过这里展开说。ForkJoinPool.commonPool() 是 JVM 全局共享的默认并行度是 Runtime.getRuntime().availableProcessors() - 1。在容器化环境中这个值取决于 JVM 是否感知到了容器配额。Java 8u191 之前JVM 读的是宿主机的核数而不是容器限额一个限制在两核的 Pod 可能看到六十四核commonPool 就开了六十三个线程上下文切换开销直接把性能拖垮。即使版本够新commonPool 仍然有三个硬伤。第一所有模块共享。你引入的第三方库、框架内部的异步实现、自己写的业务代码都用同一个池子。任何一个地方发生阻塞都会影响到其他所有使用者。这种故障串联在生产环境里是非常危险的。第二没有队列控制。commonPool 内部的任务队列是无界的任务提交速度超过消费速度时任务会无限堆积直到内存被撑爆。你没法给它设一个上限来触发快速失败。第三不可监控。你拿不到它的活跃线程数、队列长度、拒绝次数出了问题只能靠 jstack 抓现场排查效率极低。所以结论很直接生产环境的业务异步一律用自建线程池把线程数、队列长度、拒绝策略、线程名前缀、监控埋点全部握在自己手里。3.2 线程数到底开多少从公式推导到场景修正线程池大小的经典公式是这样的线程数 CPU 核数 × 目标 CPU 利用率 × (1 等待时间 / 计算时间)这个公式的来源是把线程分成两种状态跑在 CPU 上的时间和等待时间I/O、锁、网络。等待占比越高需要越多的线程来把 CPU 喂饱。举个具体的计算过程。假设部署在四核机器上目标 CPU 利用率定 0.8避免跑满导致排队。一次下游调用中序列化请求、拼装参数、解析响应的纯计算部分加起来约 5ms而网络往返加下游处理约 95ms。那么等待时间与计算时间的比值就是 95 / 5 19。代入公式4 × 0.8 × (1 19) 64。也就是说理论上需要六十四个线程才能让 CPU 保持 80% 的利用率。这个数字看起来很大但在纯 I/O 等待的场景下是合理的——线程大部分时间都在等网络CPU 是空闲的。不过公式给的是理论上限实际落地的数值还要做几轮修正。第一轮看下游的承载能力。六十四个线程意味着对下游可能有六十四个并发请求如果下游是个承载能力只有三十并发的老服务你开六十四就是把它打死然后大家一起超时。这时候必须按下游容量倒推宁可让上游排队。第二轮看总资源。每个线程大约占 1MB 栈空间默认 -Xss六十四线程就是 64MB加上堆内存得算一下容器的内存限制够不够。另外还要考虑你的服务同时对外提供多少并发接受的请求数乘以每个请求的线程占用不能超过池子容量太多否则大量请求会在队列里排队延迟反而更高。第三轮做压测确认。公式算出来的永远只是起点一定要在生产同规格的环境上压一轮看几个关键指标线程池活跃数、队列长度、任务平均等待时间、下游错误率。逐步调整线程数找到吞吐量不再上升、但延迟开始恶化的那个拐点往回收 10% 到 20% 就是比较合适的配置。我实际项目中用得比较多的一个经验值是纯 I/O 型的远程调用线程池线程数取 CPU 核数的 8 到 16 倍队列长度取线程数的 2 到 4 倍剩下的用拒绝策略兜底。这个范围不是拍脑袋是在几个不同量级的服务上反复压测收敛出来的。3.3 线程池隔离与监控埋点线程池隔离的原则很简单不同性质的任务不能共用一个池子。我通常会按这个维度切分远程调用类HTTP、RPC一个池子本地缓存访问如果用了阻塞客户端再开一个小池子CPU 密集的计算一个池子定时任务一个池子。这样做的价值在于故障隔离——某一个下游变慢只会打满它对应的那个池子不会影响其他链路。线程命名前缀是低成本的排查利器。给每个池子的线程起一个有意义的前缀比如 detail-rpc-、report-calc-线上 dump 线程栈的时候一眼就能看出是谁在忙、谁堵住了。这个习惯建议从项目第一天就养成后期补非常痛苦。监控埋点至少要覆盖这几个指标用 Micrometer 或你项目里现成的监控库都能做指标采集方式告警阈值建议活跃线程数ThreadPoolExecutor.getActiveCount()持续超过核心线程数的 80%队列长度getQueue().size()超过队列容量的 60%已完成任务数getCompletedTaskCount()增速骤降需要关注拒绝任务数自定义 RejectedExecutionHandler 计数大于 0 即告警任务等待时间提交时打时间戳执行时算差值P99 超过 100ms自定义拒绝策略的地方要多说一句。ThreadPoolExecutor 内置的四种策略里AbortPolicy 会抛 RejectedExecutionExceptionCallerRunsPolicy 会让提交任务的线程自己执行也就是把压力传导回上游。这两者在异步编排里的表现完全不同。用 AbortPolicy异常会沿着 CompletableFuture 的链路传播最后被 exceptionally 捕获你能明确知道发生了拒绝。用 CallerRunsPolicy任务被主线程执行看似更平滑但如果主线程本身就是 Tomcat 的工作线程那就等于把 Web 容器的线程也拖进了这个池子的逻辑里容易引发连锁的线程耗尽。我个人的选择是核心业务链路用 AbortPolicy 加明确的降级逻辑非核心的、可以容忍延迟的用 CallerRunsPolicy。同时在拒绝处理器里打一条带线程池名字和当前状态的日志这是事后定位问题的关键证据。4. 异步链路上的异常与超时4.1 异常在链路上是怎么走的CompletableFuture 的异常传播机制和同步代码完全不一样这也是它最容易让人困惑的地方。核心规则是一旦某个阶段以异常完成后续所有不带恢复能力的阶段都会直接跳过异常一路向后传播直到遇到一个能处理它的方法。具体来说如果链路是 supplyAsync(...).thenApply(...).thenAccept(...)而 supplyAsync 里抛了异常那么 thenApply 和 thenAccept 都不会执行最终这个 CompletableFuture 是以异常状态结束的。这时候如果你调用 join()会抛 CompletionException调用 get()会抛 ExecutionException两者的原始异常都藏在 cause 里。这里有一个非常致命的坑如果你不调用 join 或 get也没有在任何地方注册异常处理那么这个异常就彻底消失了。没有任何日志没有任何告警代码静悄悄地什么都不做。这是异步编程里最难查的一类问题——业务方说这个功能没生效你看代码逻辑完全正确最后发现是某一步抛了异常没人管。所以我的硬性要求是每一条 CompletableFuture 链路的末端必须有明确的异常处理而且异常处理里必须有日志。哪怕只是 .exceptionally(ex - { log.error(xxx failed, ex); return null; }) 这么简单也比什么都没有强得多。4.2 exceptionally、handle、whenComplete 三个方法怎么选这三个方法都能接触异常但用途不同选错了要么编译不过要么行为不符合预期。exceptionally(FunctionThrowable, T)只在发生异常时触发返回一个同类型的兜底值把异常吞掉让链路恢复成正常完成状态。适合做降级比如远程调用失败时返回一个默认的空对象。handle(BiFunctionT, Throwable, R)无论成功还是失败都会触发接收结果和异常两个参数返回一个新的值。它既能处理异常也能处理正常结果而且可以改变返回类型。适合不管成不成功我都要重新组织一下返回结构的场景。whenComplete(BiConsumerT, Throwable)无论成功还是失败都会触发但不改变结果。它像是 finally 块用来做收尾动作比如释放资源、记录耗时、上报监控。它返回的 CompletableFuture 仍然保持原来的结果或异常状态异常会继续往后传。一个常见的误用是用 exceptionally 做了降级然后以为链路后面的 whenComplete 还能感知到原始异常——实际上感知不到了因为异常已经被前面吞掉了。所以在设计恢复顺序的时候恢复操作要放在收尾操作之后或者干脆用 handle 一次性把两者都处理掉。CompletableFutureDetailVO safe rawFuture .exceptionally(ex - { log.error(detail assemble failed, orderId{}, orderId, ex); metrics.counter(detail.fail).increment(); return DetailVO.empty(); // 降级兜底 }) .whenComplete((vo, ex) - { // 走到这里时 ex 一定是 null因为上面已经恢复过了 log.info(detail assembled, cost{}ms, System.currentTimeMillis() - start); });4.3 超时控制orTimeout 与 completeOnTimeout 的差别没有超时的异步链路等于没有链路。一个卡住的下游会顺着链路一路把线程池占满然后整个服务雪崩。CompletableFuture 在 Java 9 之后提供了两个原生超时方法差别值得说清楚。orTimeout(long timeout, TimeUnit unit)到达超时时间还没完成就让这个 future 以 TimeoutException 异常完成。链路会因为异常往下走最终触发你的 exceptionally 降级逻辑。completeOnTimeout(T value, long timeout, TimeUnit unit)到达超时时间还没完成就用给定的默认值让 future 正常完成。它不会产生异常链路继续往下走只是拿到的值是兜底值。选择依据很简单如果你的业务允许超时了就用默认值顶上用 completeOnTimeout 会让代码更干净如果需要明确区分正常返回和超时降级这两种情况用 orTimeout 配合 exceptionally。// Java 9 CompletableFuturePrice priceF CompletableFuture .supplyAsync(() - priceService.get(sku), ioPool) .orTimeout(200, TimeUnit.MILLISECONDS) .exceptionally(ex - { log.warn(price timeout, fallback to default, sku{}, sku); return Price.defaultValue(); });如果你还在 Java 8没有这两个方法得自己实现。标准做法是用一个定时器在超时后调用 completeExceptionally或者用 ScheduledExecutorService 配合一个中间 CompletableFuture。这个自己实现的版本要注意清理定时任务否则大量超时的任务会堆积成内存泄漏。注意orTimeout 只是让 CompletableFuture 对象进入完成状态它并不会中断底层正在执行的任务。下面那个还在阻塞的 HTTP 调用依然会继续跑直到它自己超时或者返回。所以真正要控制资源占用还是得给底层的 HTTP 客户端、RPC 框架设置连接超时和读取超时双重保险才稳。5. 实战订单详情页并行聚合改造5.1 改造前的耗时基线这个案例来自一个真实的订单详情接口。改造前它是一个典型的串行实现先查订单主表拿基本信息约 60ms再根据订单里的用户 ID 查用户信息约 50ms再查商品信息约 70ms再查物流状态约 120ms最后查可用的营销活动约 90ms。五个调用串起来加上自身的处理逻辑P50 大约 400msP95 到了 520ms。问题在于这五个调用里物流和营销这两个是比较慢的下游而且调用量一大的时候波动很明显。它们在链路的末端前面的调用把最坏延迟都累积进去了导致尾延迟特别难看。这是串行聚合的典型症状平均值还行P99 一塌糊涂。优化的空间也很清楚这五个调用里用户信息和商品信息只依赖订单里的 ID和订单主表的查询结果没有强依赖关系完全可以并行。物流和营销也是同理。真正有依赖的只有先查订单拿到 ID这一步其他四步都能并行。5.2 依赖关系拆解与编排设计拆解之后整个流程分成两段。第一段是订单主表查询必须最先执行因为后面的用户 ID、商品 ID 都从它里面取。这一段没法并行。第二段是四个可以并行的查询用户信息、商品信息、物流状态、营销活动。这四个都拿到订单 ID 之后就可以同时发起。第二段之后还有一个组装阶段需要四个结果都到齐才能拼装 DetailVO。这一步必须等全部完成。对应的编排结构就是先 supplyAsync 查订单然后 thenCompose 进入并行阶段并行阶段用四个 supplyAsync 发起四个任务用 allOf 等待全部完成然后 thenApply 组装。还有一个额外的设计考虑物流和营销这两个慢且非核心的调用应该做超时降级。物流查不到就显示暂无物流信息营销查不到就不展示营销模块不能因为它们拖慢整个页面。用户信息和商品信息是核心超时的话应该让整个请求失败返回错误提示。5.3 代码落地聚合、降级、超时public DetailVO getOrderDetail(String orderId) { long start System.currentTimeMillis(); CompletableFutureDetailVO chain CompletableFuture // 第一段查订单主表后续所有调用都依赖它 .supplyAsync(() - orderService.getOrder(orderId), rpcPool) // 进入并行阶段 .thenCompose(order - { String userId order.getUserId(); String skuId order.getSkuId(); // 核心数据超时则整体失败 CompletableFutureUser userF CompletableFuture .supplyAsync(() - userService.get(userId), rpcPool) .orTimeout(300, TimeUnit.MILLISECONDS); CompletableFutureItem itemF CompletableFuture .supplyAsync(() - itemService.get(skuId), rpcPool) .orTimeout(300, TimeUnit.MILLISECONDS); // 非核心数据超时降级为空不影响整体 CompletableFutureLogistics logisticsF CompletableFuture .supplyAsync(() - logisticsService.query(orderId), rpcPool) .orTimeout(250, TimeUnit.MILLISECONDS) .exceptionally(ex - { log.warn(logistics timeout, degrade, orderId{}, orderId); return Logistics.empty(); }); CompletableFuturePromotion promoF CompletableFuture .supplyAsync(() - promoService.query(userId, skuId), rpcPool) .orTimeout(250, TimeUnit.MILLISECONDS) .exceptionally(ex - { log.warn(promotion timeout, degrade, orderId{}, orderId); return Promotion.empty(); }); // 等四个都完成再组装 return CompletableFuture .allOf(userF, itemF, logisticsF, promoF) .thenApply(v - { DetailVO vo new DetailVO(); vo.setOrder(order); vo.setUser(userF.join()); vo.setItem(itemF.join()); vo.setLogistics(logisticsF.join()); vo.setPromotion(promoF.join()); return vo; }); }) .exceptionally(ex - { log.error(order detail failed, orderId{}, orderId, ex); throw new BizException(订单详情加载失败请稍后重试); }) .whenComplete((vo, ex) - { long cost System.currentTimeMillis() - start; metrics.timer(order.detail.cost).record(cost, TimeUnit.MILLISECONDS); log.info(order detail done, orderId{}, cost{}ms, orderId, cost); }); return chain.join(); }有几处细节值得单独说明。第一rpcPool 是我专门为这个接口配的线程池核心线程数 16最大 32队列 128线程名前缀 order-detail-。这个数字是根据前面的公式和几轮压测收敛出来的不是随便填的。第二超时时间的选择也有依据。物流和营销的日常 P99 大约是 180ms 和 150ms我给的 250ms 留了约 40% 的余量既能覆盖大部分正常情况又能在下游劣化的时候及时切走。而核心数据的 300ms 是基于订单主表 P99 约 200ms 估算的留出更宽一点的空间。第三最外层的 exceptionally 重新抛了一个业务异常。这一点很重要——如果不抛接口会返回一个空的 DetailVO前端会展示成一个空白页用户完全不知道发生了什么。抛出明确异常让上层统一处理成友好的提示是更好的体验。5.4 压测结果对比在同规格的预发环境用同样的流量模型500 QPS持续 10 分钟压测结果如下指标串行版本并行版本变化P50 耗时402ms128ms-68%P95 耗时520ms176ms-66%P99 耗时890ms310ms-65%平均 CPU 使用率22%26%4pt线程池活跃数峰值821—下游错误率0.3%0.4%基本持平P99 从 890ms 降到 310ms 是收益最大的一块因为它把原来串行累加的最坏延迟给消掉了。CPU 只涨了四个百分点因为这个接口绝大部分时间都在等 I/O。下游错误率基本没变说明并行并没有给下游增加过大的压力——这一点很关键如果并行之后错误率明显上升就得回头检查线程数和下游容量的匹配关系。还有一个不在表格里但很明显的收益代码结构变得清晰了。串行版本里那五个调用各自的 try-catch、各自的判空、各自的默认值处理散落在几十行代码里。改成编排之后每个调用的失败策略都和它自己写在一起读代码的时候不需要上下跳。6. 常见问题与排查技巧实录6.1 常见问题速查表现象大概率原因排查手段解决方向链路完全不执行无日志无异常某阶段抛异常且无处理在链路末端加 whenComplete 打日志补 exceptionally每段链路必须有异常出口接口偶发超时但下游监控正常使用了 commonPool 被其他任务占满检查是否漏传 Executor 参数所有异步任务显式指定业务线程池线程池活跃数长期打满下游变慢或有阻塞调用打线程栈看线程状态设超时 降级 按下游容量调线程数join() 长时间不返回前置任务被阻塞或死锁jstack 找线程栈看谁在等谁加超时检查是否有嵌套等待任务结果丢失、数据不完整allOf 里有任务异常后其他任务被忽略检查每个子任务的异常处理子任务单独做降级别让一个失败影响整体日志里 TraceId 是空的ThreadLocal 没有跨线程传递检查 MDC、TraceContext 的传递用装饰器包装线程池或手动透传上下文6.2 线程上下文丢失一个几乎人人都会踩的坑ThreadLocal 是和线程绑定的这个大家都知道。但很多人写异步代码的时候会下意识忘记主线程里 set 进去的值新线程里读不到。最典型的表现是日志里的 TraceId 变成空的链路追踪断掉出问题的时候根本串不起一次请求的完整调用链。另一类表现是用户上下文丢失异步任务里拿不到当前登录用户权限校验直接失败。解决思路有两种。一种是手动透传在提交任务之前把需要的值从 ThreadLocal 里取出来作为方法参数传到异步任务里在任务开头重新 set 进去。这种方式最直白但代码里会到处都是这种透传样板维护起来烦。另一种是包装线程池。写一个 ThreadPoolExecutor 的子类重写 execute 方法在提交任务的时候捕获当前线程的上下文快照在任务真正执行之前恢复执行完之后清理。阿里开源的 TransmittableThreadLocalTTL就是干这个的它提供了 TtlExecutors 工具类可以直接包装已有的线程池改动成本很低。// 用 TTL 包装线程池上下文自动透传 Executor wrapped TtlExecutors.getTtlExecutor(rpcPool); CompletableFuture.supplyAsync(() - { // 这里能读到父线程 set 进去的 TraceId log.info(async task traceId{}, MDC.get(traceId)); return remoteCall(); }, wrapped);提示包装线程池会带来额外的开销主要体现在每次任务提交时的上下文拷贝。如果你的链路非常短、QPS 很高要评估一下这部分开销。实测下来在常规业务量级下影响很小通常在 1% 以内。6.3 事务边界与异步混用的坑数据库事务是基于 ThreadLocal 绑定连接实现的异步线程拿不到主线程的事务上下文。这意味着在 Transactional 方法里发起一个 supplyAsync异步任务跑的是一套独立的事务甚至可能是独立的数据源连接写进去的数据主线程完全看不到。更糟糕的情况是主线程事务回滚了但异步线程里执行的那部分数据库操作已经提交了。这时候数据就处于一种不一致的状态而且非常难发现因为两边都没有报错。我的处理原则是异步任务里不做写操作。所有的写库、更新状态、发消息这类动作要么放在主线程的事务内同步执行要么在事务提交之后再通过事件机制触发异步处理。如果确实需要在异步任务里写那就让它用一个独立的事务并且明确这是补偿逻辑不是主流程的一部分。读操作倒是可以在异步任务里放心做因为它不涉及事务边界的问题。前面那个订单详情页的案例就是全读操作所以适合并行。6.4 任务堆积与线程池打满的排查路径线上出现响应变慢、超时增多的时候如果怀疑是线程池的问题可以按这个顺序排查。第一步看监控面板上的活跃线程数和队列长度。活跃数打满且队列持续增长基本可以确认是任务积压。第二步抓线程栈。用 jstack 或者 arthas 的 thread 命令关注你那个线程池前缀的线程都在干什么。如果大量线程停在 socketRead、httpClient.execute、future.get 这类方法上说明是下游变慢导致的阻塞。第三步确认是哪个下游。看线程栈里的调用链或者结合上游的依赖监控找出响应时间恶化的那个服务。这时候可以临时做两件事调低该下游的超时时间让它快速失败走降级或者临时扩容线程池先扛住流量。但扩容只是权宜之计如果下游本身扛不住扩线程池只会把压力传过去问题从上游转移到下游。第四步复盘根因做长期修复。要么是线程数配置不合理要么是缺少超时和降级要么是某个下游的容量规划出了问题。这三类问题对应的修复手段完全不同别混在一起治。还有一个容易被忽略的现象如果拒绝策略用的是 CallerRunsPolicy线程池满的时候任务会被提交线程自己执行表现为提交任务的接口变慢。这时候你去看被提交的那个线程池反而看不出什么异常因为压力被转移到调用方了。排查的时候要注意这个陷阱看看是不是 Tomcat 的工作线程在跑本应该在线程池里跑的任务。7. 一些个人取舍经验7.1 join 和 get 的选择join 抛的是 CompletionExceptionget 抛的是 ExecutionException。前者是非受检异常后者是受检异常。这个差别看起来很小但在链式代码里影响很大。如果链路上到处都是 get你就得为每个调用写 try-catch或者让方法签名抛出检查异常代码会变得很难看。所以写编排的时候我一律用 join只有在最外层需要把异常转换成业务异常的时候才用 get 或者统一 catch。还有一个性能上的细微差别在同一个 CompletableFuture 上多次调用 get 或者 join它们的开销是可以忽略的因为结果已经缓存了。所以不用担心在 allOf 的回调里连续 join 五个 future 会有性能问题。7.2 不要在 thenApply 里做阻塞调用这一条我踩过真实的坑。有一段代码是这样的supplyAsync(fetch, pool).thenApply(this::enrich)enrich 里做了一次远程调用没有加 Async也没有加超时。结果就是fetch 所在的线程执行完 fetch 之后接着执行 enrich一阻塞就是几百毫秒。在 QPS 上来之后线程池里的线程全被 enrich 占住整个池子彻底瘫掉。现在的规矩是thenApply 里只放纯内存操作任何涉及 I/O 的动作要么用 thenApplyAsync 显式切换线程池要么用 thenCompose 包一层新的 supplyAsync。判断标准就是问自己一句这段代码会不会等待会不会有不确定的耗时只要答案是可能会就切出去。7.3 保持编排的可读性比追求极致并行更重要我见过为了把五个调用全塞进一条链里、结果写成三十行嵌套 lambda 的代码也见过一个方法里拼接了七八种 API、后来的人根本不敢改的代码。异步编排的可读性和它的性能收益一样重要因为代码是要被维护好几年的。我的做法是控制单个链路的长度超过五六个阶段就拆方法。把取数据和组装数据分开取数据的方法返回一个包含若干个 CompletableFuture 的小容器对象组装的方法只负责 join 和拼接。这样两部分的职责清晰测试也好写。另外给每个 CompletableFuture 变量起一个有业务含义的名字比写成 future1、future2、future3 要重要得多半年后回头看名字就是最好的文档。最后再分享一个小技巧。调试异步链路的时候我习惯在每个关键阶段后面加一个 thenApply 打点打印当前阶段名和时间戳。在预发环境打开生产环境通过开关控制。这东西在排查到底卡在哪一步的时候非常有用比看线程栈直观得多。而且它的成本很低几行代码的事。还有一个我一直在坚持的习惯任何一条 CompletableFuture 链路的末端必须有三个东西——异常处理、日志、耗时统计。少了任何一个等线上出问题的时候你就得靠猜。这个习惯刚开始写的时候会觉得啰嗦写顺了之后会发现它是你半夜能睡好觉的底气。
网站建设高端定制企业官网