ARTICLE DETAIL

资讯详情

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

SpringBoot高并发接口优化:CompletableFuture异步编排与线程池实战

SpringBoot高并发接口优化:CompletableFuture异步编排与线程池实战 做高并发服务端开发这几年我最大的一个感触是很多接口慢不是单次调用慢而是十几个串行调用叠在一起慢。SpringBoot项目里最常见的就是一个Controller里按顺序查库存、查价格、查会员、查优惠、发通知、写日志一个接口几百毫秒就没了。后来我全面用CompletableFuture 自定义线程池做异步编排把接口P99从400多毫秒压到了150毫秒以内。这篇文章就把我从原理到落地、再到线上踩坑的完整经验写出来包括线程池参数怎么定、CompletableFuture的核心API怎么用才不会写出嵌套地狱以及几个真实翻车现场。如果你正在用SpringBoot做高并发接口或者你的接口里一堆调用之间有并行关系、有一层一层的依赖关系这篇内容可以直接参考。我不讲虚的全部给可落地的代码、参数和排查思路。1. 从同步阻塞到异步回拨为什么你的业务需要编排1.1 一个典型的串行地狱接口先看一段最常见的同步代码Service RequiredArgsConstructor public class OrderDetailService { private final OrderDao orderDao; private final UserService userService; private final CouponService couponService; private final StockService stockService; private final RecommendService recommendService; public OrderDetailVO getDetail(String orderId) { OrderDO order orderDao.selectById(orderId); // 20ms UserInfo user userService.getUserById(order.getUserId()); // 30ms CouponVO coupon couponService.calcCoupon(orderId); // 50ms StockVO stock stockService.queryStock(order.getSkuId()); // 30ms ListProductVO recommend recommendService.recommend(order); // 80ms return new OrderDetailVO(order, user, coupon, stock, recommend); } }这段代码的逻辑很清晰用户点进订单详情后端依次去查数据库、查用户服务、算优惠、查库存、拿推荐商品。假设每次调用耗时分别是20ms、30ms、50ms、30ms、80ms接口总耗时就是它们的和210ms再加上网络开销、序列化开销实际可能300ms以上。在高并发的场景下这个串行地狱有两个恶果用户端响应慢体验差App的订单详情页经常转圈。Tomcat的工作线程被长期占用。假设QPS是500平均接口耗时300ms那么活跃线程数就是500×0.3150个。线程池一旦被占满新的请求只能排队Tomcat线程上堆积的请求越来越多雪崩就是这么来的。所以要做异步编排第一个动机不是炫技而是把多个相互独立的远程调用从加法变成取最大值。本来300ms如果能并行理论上只需要最慢的那个80ms再加上最后组装的时间。1.2 不是所有异步都叫编排很多同学一听说异步就随手new Thread(() - xxxService.send())或者把方法丢进线程池就不管了。这里有个概念要分清楚异步执行和异步编排不是一回事。异步执行把任务扔给另一个线程去跑主线程不等结果。典型的就是发完订单消息就返回MQ慢慢消费。异步编排主线程发起多个任务但还需要等待这些任务的结果并且根据结果继续做下一步计算比如等并行结果全回来之后合并成一个VO返回给前端。CompletableFuture解决的就是异步编排这件事。它提供了一整套组合工具supplyAsync发起异步任务thenApply处理上一步的结果thenCombine合并两个任务的结果allOf等待一堆任务全部完成exceptionally做异常兜底。我习惯打一个比方线程池是整个工厂的施工队CompletableFuture是这个工地的总指挥。施工队负责干活总指挥负责说这两个活可以同时干那个活必须等这个活完了再干有工序失败了换条路接着干。没有总指挥你只知道把活扔下去不知道什么时候能干完、哪个干失败了。1.3 适合编排的业务画像我并不是建议把每个接口都改成异步编排。异步比同步多一层复杂性有多线程就多线程安全的问题。根据经验适合用CompletableFuture编排的业务通常符合下面几个特征中的至少两条多个下游调用之间没有数据依赖可以并行。比如订单接口里查库存、查优惠、查推荐商品之间互不影响。存在明显的先做A拿到结果后再做B的依赖关系而且A和B之间还可能插着其他并行任务。比如先查用户信息再根据用户等级去算优惠。单个下游服务不太稳定时不时慢一次需要做超时和降级。接口的响应时长对这个业务KPI很重要比如App首屏接口、下单链路。反过来如果这一步操作涉及强事务、强一致性比如扣库存和加积分必须同生共死那我不建议拆进异步编排。异步拆开之后事务边界就不好控制了我后面会专门讲这个坑。2. 线程池的线程吃饱了异步还不如同步快2.1 默认ForkJoinPool是公共泳池线上别裸用CompletableFuture有个很方便的设计如果你不传线程池它默认会用ForkJoinPool.commonPool()来执行任务。开发环境跑一下确实挺爽但线上千万不要裸用。原因有三个commonPool是全局共享的线程数默认是CPU核数 - 1。假设机器是8核它就只有7个工作线程。在高并发下这些线程一旦被IO操作阻塞后续所有依赖commonPool的异步任务全部排队。不只是你的业务在用这个池子JDK里的其他并行流、一些第三方库可能也在用。一个业务线的慢任务可能把另一个不相干链路的异步任务全部拖死。它的线程数跟机器CPU挂钩不适合IO密集型任务。依赖的是下游HTTP、RPC线程大部分时间是阻塞等待下游返回的7个线程根本不够用。我在线上遇到过真实事故一个埋点上报任务用了commonPool结果高峰期埋点服务响应变慢线程全被堵在HTTP调用上最后导致其他使用commonPool的并行流代码也一起变慢。这个案例放在文章第5部分细讲。正确的做法是只要用了CompletableFuture就显式传入一个自己可控的线程池。2.2 线程池参数怎么定从QPS反推核心线程数线程池参数不是拍脑袋定的。我的思路是从目标QPS和平均响应时间反推并发度。先看一个公式并发数 ≈ QPS × 平均响应时间秒。假设接口目标QPS是500每个异步子任务平均响应时间是200ms那么任何一个时间点同时在执行的子任务数量大概是500 × 0.2 100这个100就是理论并发任务数。所以核心线程数可以设为100左右最大线程数可以适当放宽到150到200前提是下游服务能扛得住。如果下游服务只有40个线程的容量的你这边开100个线程打过去就是把下游打挂这个参数必须跟DBA、下游服务负责人对齐。队列长度我建议用有界队列而不是Executors那种无界链表。无界队列在流量突增时会把任务无限堆积内存被打高而且线程数永远不扩张延迟越来越大。用有界队列的好处是满的时候能触发扩容线程或者触发拒绝策略我们至少能感知到超载。这里还有一个特别值得注意的点ThreadPoolExecutor的默认调度策略是先入队后扩线程。也就是说核心线程满了之后新任务先往队列里放队列满了才创建非核心线程。在突发流量下你希望的是多开线程干活而不是让请求堆在队列里排队。所以我一般自定义RejectedExecutionHandler在队列满了且线程数未达上限时直接创建线程执行如果线程数也满了再走CallerRunsPolicy让调用者线程去执行起到天然背压的效果。2.3 线程池隔离与命名别用new Thread Executors先列几个不要做的事情不要用Executors.newFixedThreadPool()底层是无界队列突发流量会让队列无限膨胀。不要用Executors.newCachedThreadPool()最大线程数是Integer.MAX_VALUE高峰期能创建出几万个线程直接把内存和下游打爆。不要在业务里new Thread()跑异步任务线程创建开销大且无法复用、无法监控。线程池隔离是什么意思我强烈建议不同业务、不同下游性质的任务用不同的线程池。比如订单详情编排用detailPool消息推送用pushPool报表聚合用reportPool。这样某个下游抖动时最多打满它自己的线程池不会连累其他编排链路。另外给线程起一个能看出来的名字非常非常重要。之前排查线上问题线程dump出来全是pool-3-thread-1这种默认名字根本不知道是哪个业务创建的。后面全部换成了带业务前缀的命名工厂public class NamedThreadFactory implements ThreadFactory { private final AtomicInteger sequence new AtomicInteger(1); private final String prefix; public NamedThreadFactory(String prefix) { this.prefix prefix; } Override public Thread newThread(Runnable r) { Thread thread new Thread(r, prefix - sequence.getAndIncrement()); thread.setDaemon(false); return thread; } }这样线程名就是order-detail-pool-1、order-detail-pool-2。一旦线上打dump一眼就能看出是哪条链路的问题。3. CompletableFuture 核心 API 的实战拆解3.1 结果加工thenApply 与 thenCompose 的差别先介绍最常用的一组串行加工APIthenApply和thenCompose。thenApply的意思是上一个阶段的结果已经拿到我把它加工成另一个值同步返回。CompletableFutureInteger future CompletableFuture .supplyAsync(() - getUserService().getUserAge(1001), detailPool) .thenApply(age - age 10);thenCompose的意思是上一个阶段的结果要用来发起另一个异步任务。它的返回值不是普通值而是一个新的CompletableFuture。CompletableFutureUserInfo future CompletableFuture .supplyAsync(() - userService.getUserById(1001), detailPool) .thenCompose(user - userService.getOrdersByUserIdAsync(user.getId()));如果这里用thenApply而不是thenCompose就会产生CompletableFutureCompletableFutureListOrderDO这种嵌套结构处理起来非常别扭。我的经验是只要下一步是发起新异步调用就用thenCompose下一步是拿到结果做计算就用thenApply。3.2 并行聚合thenCombine 与 allOf 谁更适合你thenCombine适合两个独立任务执行完把俩结果合并成一个CompletableFuturePriceResult future CompletableFuture.supplyAsync(() - priceService.calcBasePrice(skuId), detailPool) .thenCombine( CompletableFuture.supplyAsync(() - couponService.calcCoupon(userId), detailPool), (basePrice, coupon) - new PriceResult(basePrice, coupon) );如果并行任务超过两个甚至是一批动态数量的任务thenCombine就不好使了。这时候用allOfListCompletableFutureObject futures new ArrayList(); futures.add(CompletableFuture.supplyAsync(() - orderDao.selectById(orderId), detailPool)); futures.add(CompletableFuture.supplyAsync(() - userService.getUserById(userId), detailPool)); futures.add(CompletableFuture.supplyAsync(() - stockService.queryStock(skuId), detailPool)); CompletableFutureVoid all CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])); all.join(); OrderDO order (OrderDO) futures.get(0).join(); UserInfo user (UserInfo) futures.get(1).join(); StockVO stock (StockVO) futures.get(2).join();这里有个细节注意allOf本身不聚合结果它只负责等所有任务完成。想要结果还得逐个join()。如果任务很多强烈不建议用futures.get(i).join()这种下标的写法一旦有人调整了列表顺序数据就错位了。更稳妥的是给每个任务单独命名或者用Map按业务key存future。再对比一下join()和get()join()不检查受检异常失败时抛CompletionException代码干净。get()抛出InterruptedException和ExecutionException需要try-catch。线上我一般用join()然后配合全局异常处理器把CompletionException里的真实原因剥出来。3.3 兜底失败exceptionally、handle 与 whenComplete 的差异这三个方法都会在任务出现异常时触发但语义差别很大选错会出现很隐蔽的问题。exceptionally只处理异常分支返回一个代替值CompletableFutureStockVO stockFuture CompletableFuture .supplyAsync(() - stockService.queryStock(skuId), detailPool) .exceptionally(ex - { log.warn(query stock failed, skuId{}, skuId, ex); return StockVO.unknown(); });handle是无论成功还是失败都会执行它能同时拿到结果和异常必须返回一个结果CompletableFutureString future CompletableFuture .supplyAsync(() - stockService.queryStock(skuId), detailPool) .handle((result, ex) - { if (ex ! null) { return fallback; } return result.getStatus(); });whenComplete更像是一个观察者回调它拿到结果和异常但不能改变计算结果future.whenComplete((result, ex) - { if (ex ! null) { metrics.recordError(stock_query); } else { metrics.recordTime(stock_query, result.getCostMs()); } });我用一句话总结想降级就exceptionally想统一收尾就whenComplete想在失败成功两种情况下做不同计算就用handle。监控上报我放在whenComplete里这个习惯帮我排查过很多次慢调用。4. 一个完整的高并发异步编排落地案例4.1 需求与接口设计讲原理很容易飘落到项目里才有意义。我拿一个真实改造过的场景举例订单详情页接口。接口路径是GET /order/detail/{orderId}前端在打开订单列表、支付回调页、订单详情页时都会调用。改造前它长这样查询订单主记录数据库约20ms根据订单用户ID查用户信息RPC约30ms计算优惠价格RPC约50ms查库存状态RPC约30ms查推荐商品列表RPC约80ms因为全部串行接口P99到过400ms。业务方要求订单详情页首屏必须在200ms以内返回数据。我们的改造目标不是让每一路调用变快而是让它们在时间轴上重叠起来。改造思路是把五路调用并发出去等全部返回后组装为VO。整体耗时的理论上限从五路之和变成最大的一路耗时 组装耗时。由于最慢的推荐商品也只有80ms所以理论上P99可以被压到150ms以内。4.2 线程池配置类怎么写才不翻车我习惯在SpringBoot里用ThreadPoolTaskExecutor因为它能被Spring容器管理可以和Async配合也方便在配置类里统一调优。Configuration public class AsyncPoolConfig { Bean(orderDetailExecutor) public ThreadPoolTaskExecutor orderDetailExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(80); executor.setMaxPoolSize(200); executor.setQueueCapacity(200); executor.setKeepAliveSeconds(60); executor.setThreadNamePrefix(order-detail-pool-); executor.setWaitForTasksToCompleteOnShutdown(true); executor.setAwaitTerminationSeconds(10); executor.setRejectedExecutionHandler(new CallerRunsPolicy()); executor.initialize(); return executor; } }几个参数解释一下corePoolSize80按前面反推的并发数100往下留了点余量因为不是所有子任务都会同时执行。maxPoolSize200允许突发流量时临时扩到200个线程但前提是下游服务能承受。queueCapacity200有界队列任务堆积超过200时触发拒绝策略。CallerRunsPolicy当线程池满时不让新请求直接失败而是让Tomcat的请求线程自己执行任务。这样接口响应会变慢但不会因为排队导致线程耗尽。对用户侧来说降级为慢请求比直接报错更能接受。setWaitForTasksToCompleteOnShutdown(true)awaitTerminationSeconds(10)应用停机时让正在跑的任务尽量完成避免服务重启时丢数据。初始化线程池后在使用时通过构造注入或者Qualifier取Bean即可。不要每次都new一个线程池那就是白配置了。4.3 用 allOf supplyAsync 编排完整接口下面是Service层的核心逻辑。我用列表来收集future然后用allOf统一等待Service RequiredArgsConstructor public class OrderDetailAssembleService { private final OrderDao orderDao; private final UserService userService; private final CouponService couponService; private final StockService stockService; private final RecommendService recommendService; Qualifier(orderDetailExecutor) private final ThreadPoolTaskExecutor orderDetailExecutor; public OrderDetailVO assemble(String orderId) { CompletableFutureOrderDO orderFuture CompletableFuture .supplyAsync(() - orderDao.selectById(orderId), orderDetailExecutor) .exceptionally(ex - { log.error(query order failed, orderId{}, orderId, ex); return null; }); CompletableFutureUserInfo userFuture orderFuture.thenCompose(order - CompletableFuture.supplyAsync(() - userService.getUserById(order.getUserId()), orderDetailExecutor) .exceptionally(ex - { log.error(query user failed, userId{}, order.getUserId(), ex); return UserInfo.unknown(); }) ); CompletableFutureCouponVO couponFuture CompletableFuture .supplyAsync(() - couponService.calcCoupon(orderId), orderDetailExecutor) .exceptionally(ex - { log.error(calc coupon failed, orderId{}, orderId, ex); return CouponVO.empty(); }); CompletableFutureStockVO stockFuture CompletableFuture .supplyAsync(() - stockService.queryStock(orderFuture.join().getSkuId()), orderDetailExecutor) .exceptionally(ex - { log.error(query stock failed, orderId{}, orderId, ex); return StockVO.unknown(); }); CompletableFutureListProductVO recommendFuture CompletableFuture .supplyAsync(() - recommendService.recommend(orderFuture.join()), orderDetailExecutor) .exceptionally(ex - { log.error(recommend failed, orderId{}, orderId, ex); return Collections.emptyList(); }); CompletableFuture.allOf(userFuture, couponFuture, stockFuture, recommendFuture).join(); return new OrderDetailVO( orderFuture.join(), userFuture.join(), couponFuture.join(), stockFuture.join(), recommendFuture.join() ); } }这里有三个实战中的关键点不要在主线程里等着拿order再决定要不要并行。我在写第一版的时候先orderFuture.join()取出OrderDO再往后编排。结果订单主记录这条同步路径变成了全链路瓶颈等于白改造。后来改成先往CompletableFuture里塞依赖它的地方再用thenCompose或者join处理。每个子任务都要有exceptionally兜底。异步任务一旦某个环节抛异常allOf会立刻把整个future标记为异常导致join()直接抛异常。所以我在每个任务内部都把异常吃掉并返回一个安全的默认值。这样单个下游挂了接口不至于直接报错。join放在allOf之后。因为allOf保证所有future都已完成此时join()不会阻塞也不会死锁。如果有人把join()写到了所有任务发起之前就会出现主线程等子线程、子线程还等在队列里的尴尬场面。上面的demo中coupon和stock跟order没有严格依赖但recommend依赖order内容。所以在recommendFuture里我直接用了orderFuture.join()因为它一定已经完成。这是用了thenCompose会更优雅但考虑到代码可读性我保留了join的方式前提是确认顺序无误。4.4 耗时实测与容量估算改造完我做了两轮压测环境和结果如下环境4核8G的测试机下游模拟接口平均耗时30到80msP99约120ms。方案P99耗时P50耗时说明未改造串行420ms380ms五路调用依次执行大量时间花在等待CompletableFuture编排不超时145ms132ms五路并行瓶颈是80ms的推荐服务CompletableFuture编排 200ms超时198ms135ms超时兜底异常场景下仍有保障做容量估算时不要只看接口QPS。我习惯于按线程池最大线程数 × 下游平均耗时来算。比如maxPoolSize200下游服务平均处理50ms那么线程池理论上每秒钟最多处理200 ÷ 0.05 4000个子任务如果一个订单详情接口并发5个子任务那么它能支撑的订单详情QPS大约是800。如果下游平均耗时涨到200ms这个数字就变成了200 ÷ 0.2 1000个子任务再除以5就是200QPS。所以线程池参数不是一劳永逸下游变慢了你的处理能力会跟着掉必须监控线程池的排队深度和活跃线程数。我还额外做了一件事给每个子任务统一加了超时控制避免某个下游服务彻底挂掉时请求线程全部挂在等待上。CompletableFutureStockVO stockFuture CompletableFuture .supplyAsync(() - stockService.queryStock(skuId), orderDetailExecutor) .orTimeout(200, TimeUnit.MILLISECONDS) .exceptionally(ex - { log.warn(stock query timeout, skuId{}, skuId, ex); return StockVO.unknown(); });orTimeout在超时后会把这个future置为异常状态配合exceptionally就能走降级返回。这个组合我强烈建议加在每个外部调用上后面第5节我会讲为什么单纯的orTimeout也救不了你。5. 线上故障复盘我的三个真实翻车现场5.1 翻车一commonPool被埋点任务拖垮这是我第一次在业务里大规模使用CompletableFuture时的教训。当时有个埋点上报功能从用户请求里采集行为数据后通过HTTP发送给数据分析平台。代码写得很简单直接用了commonPoolCompletableFuture.runAsync(() - trackingClient.send(event));因为埋点对延时不敏感当时觉得反正丢到线程池里响应快就行。结果有一天线上监控全面告警核心接口P99从80ms涨到了1.2秒而且Tomcat线程池经常打满。排查过程很有意思。先看CPU不高看GC也不异常。后来dump线程发现大量线程停留在sun.misc.Unsafe.park上继续追踪是阻塞在HTTP客户端的连接池获取上。进一步看这些线程的线程名全是ForkJoinPool.commonPool-worker-*。原因一下就清楚了机器是8核commonPool只有7个线程。埋点上报的下游数据分析平台网关抖动HTTP连接迟迟不返回7个线程全被阻塞在IO等待上。其他所有使用commonPool的业务包括订单详情里的并行流、其他同事的CompletableFuture全部在排队等待这7个线程。修复方案是埋点上报改成独立的线程池跟业务编排链路彻底隔离。从那以后我在代码里搜索了所有没传第二个参数的supplyAsync/runAsync全部显式传入业务线程池。同时也养成了一个习惯任何CompletableFuture调用不写线程池比写线程池更容易出事故。5.2 翻车二数据库连接池被异步任务占满第二次翻车是异步编排里跑了数据库操作。当时做的是一个后台批量任务需要从一批订单里读取数据加工后批量写入更新表。我用了一个核心线程数80的线程池提交800个任务每个任务内部走MyBatis查询和更新。上线后数据库连接池HikariCP的活跃连接数直接打满数据库端出现大量连接等待紧接着所有读写库的同步业务也都开始报connection is not available。为什么因为同步请求下一条Tomcat线程同一时间只占一个数据库连接。异步编排后80个甚至更多线程同时去竞争数据库连接HikariCP默认最大连接数只有10的话剩下的线程全部阻塞在getConnection()上。而每个被占用的连接执行SQL再快也要几十毫秒后面排队的线程越多整体越慢慢慢就形成了雪崩。这个事给我两条教训任务依赖数据库连接池时一定要评估连接池容量。如果你开了80个并发线程连接池至少也要有80个连接否则并发越高等待越严重。数据库密集型操作不应该塞进无界线程池。更稳妥的方式是控制并发度比如用信号量限制同时执行的任务数或者干脆用有界队列把积压控制在可控范围。后来我把这类任务的线程数降到20HikariCP最大连接数调到50同时在每个任务里显著缩短了单个事务的时间问题才算真正解决。5.3 翻车三orTimeout超时了底层任务却还在跑第三次翻车是我给每个CompletableFuture加了orTimeout(2000)之后发现接口是降级返回了但下游服务仍然时不时被打到报警。一开始我以为是下游自己的问题后来一看连接数发现来自我这边的大量请求一直悬挂在下游没有释放。这里必须把CompletableFuture的机制说清楚orTimeout在超时后只是把这个future标记为异常让调用链路上的join()或exceptionally立刻感知到。它并不会有真正去取消正在执行的底层任务。也就是说HTTP请求已经发出去了底层线程仍然阻塞在等待响应上该占的线程、该占的连接一个都没少。这个问题的本质是超时控制只解决主线程别等太久不解决底层资源赶紧释放。对下游服务来说超时的请求依然在打它。我的修复方案分两层第一层在HttpClient层面设置全局连接和读取超时让底层HTTP调用本身就有刹车而不是依赖CompletableFuture的orTimeout兜底。第二层在提交任务前用信号量限制并发。相当于给这个下游的请求总量加了水龙头超过并发上限的请求直接走降级而不是全部压到下游private final Semaphore stockSemaphore new Semaphore(30); public CompletableFutureStockVO queryStockWithLimit(String skuId) { if (!stockSemaphore.tryAcquire()) { return CompletableFuture.completedFuture(StockVO.unknown()); } return CompletableFuture .supplyAsync(() - { try { return stockService.queryStock(skuId); } finally { stockSemaphore.release(); } }, orderDetailExecutor) .orTimeout(200, TimeUnit.MILLISECONDS) .exceptionally(ex - StockVO.unknown()); }这一层信号量本质上是对下游的关键保护。我现在做任何一个异步编排依赖外部服务时都会问一句这个下游的最大承受并发是多少如果答不上来就先按30、50这类保守值加信号量后续根据监控逐步调整。最后再分享一个我的实际体会异步编排不是把代码改成CompletableFuture就完事了。线程池参数、超时策略、降级兜底、并发控制、线程名规范、监控指标六样东西缺一不可。我建议你在SpringBoot项目里把ThreadPoolTaskExecutor的队列深度、活跃线程数、拒绝任务数接到监控系统里每次上线后先看这些指标再决定要不要调大并发。这套组合跑下来你的接口会比我改造前的订单详情接口快很多但不会有我之前踩过的这些坑。
返回列表