拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

SpringBoot高并发异步编排:CompletableFuture与线程池实战指南

SpringBoot高并发异步编排:CompletableFuture与线程池实战指南

做后端这几年,被高并发折磨得最厉害的往往不是数据库,而是那种“一个接口要串行调好几个服务”的场景。一个典型的 SpringBoot 项目里,订单详情要查订单、查库存、查优惠券、查用户,每个服务稳定 50ms,串起来就是 200ms。压测一到 1000 QPS,Tomcat 线程全被这 200ms 的等待占满,接口大面积超时,CPU 却闲得很。CompletableFuture 加线程池做高并发异步编排,就是干这个用的:把相互独立的调用拆开并行执行,再用编排原语把结果汇拢,把 RT 从“多个 RTT 之和”压到“最慢的那个 RTT”。

这篇文章我不打算讲太多理论,重点放在 SpringBoot 里这套方案怎么落地,线程池怎么配才不会被拖垮,CompletableFuture 的常用编排怎么用最短代码写出来,以及几个我实际踩过且线上真实发生的坑。适合正在写高并发接口、对 Future 理解还不深、想把异步编排用明白的 Java 工程师。看完可以直接照着抄。

1. 为什么要做异步编排:这笔账其实很好算

1.1 串行等待是怎么拖垮接口的

先算一笔直白的账。假设一个接口依赖四个上游服务,每个服务平均返回时间都是 50ms。串行调用时,接口 RT 约等于 200ms,如果再算上网络抖动、超时重试,实际 P99 可能到 400ms 往上。

你看这 4 个线程在干嘛:一个请求进来,线程 A 占着等用户服务,线程 B 占着等订单服务,每个线程大部分时间都在阻塞。Tomcat 或你自己的业务线程池一共就那么多线程,全被这种无意义的等待耗光,新的请求只能排队。这就是高并发下最常见的瓶颈:并发上不去不是 CPU 不够,而是线程都被 IO 等待堵住了。

那并行之后呢?四个调用同时发出去,接口 RT 从“四个 RTT 之和”变成“最慢的那个 RTT”,理想情况下就是 50ms 多一点点。同样 200 个线程,原来最多撑 1000 QPS(这里只是粗略估算),现在能撑 4000 QPS,直接是一个量级的提升。

用生活里的例子更好理解:串行就像你点外卖时,只有一家餐厅、一个厨师,菜一道一道做;并行就像同时让四家餐厅分别做四个菜,你只需要等最慢的那家。CompletableFuture 没发明任何新东西,它只是把“同时让四家餐厅开工”这件事表达得足够顺手。

1.2 为什么选 CompletableFuture 而不是自己写多线程

很多老项目用 ExecutorService 加 Future 也能做并行,为什么还要换 CompletableFuture?因为老的 Future.get() 是阻塞的,而且表达不了“这个任务依赖上一个任务的结果”这种流程。

举个例子,你要查完订单才能拿订单里的用户 ID,再拿用户 ID 去查用户信息,这就是典型的串行依赖。用老 Future 写,你得先 submit 一个任务,get 到订单,再 submit 一个任务,get 到用户。中间只要有任何异常、超时,代码就变得乱七八糟。更别说“查完订单和查完库存之后,再把两个结果合起来做一件事”这种场景,手写起来全是样板代码。

CompletableFuture 的核心价值是编排,不是并发。它把任务之间的依赖关系声明出来:A 完成后把结果交给 B、A 和 B 都完成后汇合到 C、多个任务里任意一个成功就返回。这套东西在 SpringBoot 生态里用起来非常顺,因为你只需要把线程池作为参数传进去,剩下的逻辑全部是声明式的。

2. 线程池配置:先建地基再盖楼

2.1 为什么不能依赖默认的 ForkJoinPool

用 CompletableFuture 不传线程池,它会走 ForkJoinPool.commonPool()。这个池听起来挺正规,但它有非常明显的短板:

  • 默认并发度是 CPU 核数减 1。也就是说,一台 8 核的机器,commonPool 的并行度只有 7。如果你的业务全是 IO 等待(调用外部服务、查 Redis、查数据库),这 7 个线程全部被阻塞,后面所有任务全在排队。
  • commonPool 是 JVM 进程级共享的。你项目里别的地方用了 parallelStream、别的第三方库也用了 commonPool,大家都在抢这 7 个线程。其中一个任务被阻塞,可能把旁边完全不相干的任务也拖死。

高并发场景下,线程池必须是隔离的、有名字的、参数可控的。页面查询用一个池、异步推送用一个池、定时任务用另一个池。各池互不影响,出问题也好排查——线程 dump 里看到biz-io-3立刻就知道这是哪条链路的线程。

2.2 线程数怎么定:先算账,再微调

线程池参数不用拍脑袋,可以按业务模型的公式推一个起点。对于 IO 密集型任务,最常用的估算思路是看“同时在途任务数”:

假设你的服务单机目标 QPS 是 200,单个任务平均耗时(RTT)是 80ms。那么同时处理中的任务数大约是 QPS × RT = 200 × 0.08 = 16 个任务。也就是说,线程池稳态需要约 16 个线程在工作。

不过流量一定有毛刺,还要留余量。实践里我一般按这个起点配:

参数值依据
corePoolSize16稳态 16 个并发任务
maxPoolSize32留 100% 余量应对突刺
queueCapacity200允许一部分任务排队等待
keepAliveTime60s闲置线程回收的合理时间
threadNamePrefixbiz-io-出问题时能准确定位线程来源
拒绝策略CallerRunsPolicy宁可让调用方线程执行,不愿丢任务

这只是起点,不是终点。真实上线前我会用压测去校核:如果线程池活跃数长期 0,就说明配大了;如果队列一直积压几百个任务,就该扩池或加机器。线程池参数不是“配一次就完了”,它应该跟随业务节奏动态调整。

2.3 阻塞队列到底该选哪种

热词榜里有人专门搜“线程池的阻塞队列选择”,说明这事坑真的很多。ThreadPoolExecutor 执行逻辑是:核心线程满 -> 任务进队列 -> 队列满 -> 创建新线程到最大线程数 -> 再满就走拒绝策略。

所以队列类型直接决定了线程池的扩缩行为:

队列类型特点适用场景隐患
SynchronousQueue不缓存任务,直接把任务交给线程纯 CPU 密集型任务队列通常为 0,容易频繁触发 maxPoolSize
LinkedBlockingQueue默认容量是 Integer.MAX_VALUE,无界对任务丢失极其敏感的场景任务无限堆积导致内存溢出,且 maxPoolSize 永远不起作用
ArrayBlockingQueue有界队列,容量可控大多数 IO 密集型业务需要配合合理的拒绝策略,否则容易丢任务
PriorityBlockingQueue有优先级排序希望重要任务先执行无法和 FIFO 保证完全一致,排查时心智负担重

我的真实建议是:IO 密集型业务用有界队列。容量不要拍脑袋,通常按“稳态并发数 16 × 5~10”来定,也就是 100~200 之间。队列太大,任务积压一个量级,用户都等到超时了你才发现;队列太小,线程频繁创建销毁或直接进拒绝策略,接口毛刺严重。

至于拒绝策略,线上我一般优先选 CallerRunsPolicy。为什么?因为它最“诚实”:线程池忙不过来时,让打过来的调用方线程自己执行任务。任务不会丢,代价是调用方线程被占用,这个信号通过监控很容易发现。AbortPolicy 适合任务绝对不能丢但可以快速失败的场景;DiscardPolicy 适合丢一些低价值任务也无所谓的场景,比如打点上报。

3. CompletableFuture 核心 API:把编排语言讲明白

3.1 提交异步任务:supplyAsync 与 runAsync

最简单的用法,是给线程池提交一个任务,异步执行完返回结果:

Executor executor = bizExecutor; // 有返回值的异步任务 CompletableFuture<Integer> countFuture = CompletableFuture.supplyAsync(() -> orderService.countToday(uid), executor); // 无返回值的异步任务 CompletableFuture<Void> logFuture = CompletableFuture.runAsync(() -> logService.record(uid), executor);

这里有一个非常关键的细节:第二个参数不传,CompletableFuture 会默认走 ForkJoinPool.commonPool()。我在代码评审时见过很多次这种问题——开发写了supplyAsync(() -> xxx)没传线程池,本地测试一切正常,到生产环境并发一高就卡死,因为默认池太弱了。

supplyAsync和runAsync的路数也不一样。前者代表“异步跑了之后我还要拿结果”,后者代表“你只管执行,我不关心返回值”。用supplyAsync和runAsync的边界,别混用,不然还要再包一层 CompletableFuture。

3.2 串行编排:thenApply、thenCompose、thenAccept

串行任务的意思是:A 完成之后,拿 A 的结果去跑 B。CompletableFuture 提供了三个非常容易混淆的方法:

  • thenApply:上一个阶段的结果进来,经过同步转换,返回一个新结果
  • thenCompose:上一个阶段的结果进来,返回一个新的 CompletableFuture,用于扁平化
  • thenAccept:上一个阶段的结果进来,消费掉,没有返回结果

先看thenCompose为什么比thenApply更合适做串行异步:

// 错误示范:嵌套 CompletableFuture CompletableFuture<CompletableFuture<User>> badFuture = orderFuture.thenApply(order -> userService.findAsync(order.getUid(), executor)); // 正确写法:扁平化串行 CompletableFuture<User> goodFuture = orderFuture.thenCompose(order -> userService.findAsync(order.getUid(), executor));

thenCompose相当于把两层包装摊平,返回的还是一个 CompletableFuture ,后续继续用thenApply、exceptionally都很自然。新手更容易犯的错是在thenApply里做耗时操作——记住,thenApply里的代码默认是同步执行的,如果里面调了外部服务,一样会阻塞当前线程。真正想把“拿到结果后再异步查别的服务”表达出来,就要用thenCompose。

thenAccept用在哪?比如异步等待任务完成之后,把结果写进缓存,不需要向上返回。注意它依然可能阻塞:如果 CompletableFuture 是在线程池里完成,thenAccept默认也在同一个线程执行,所以里面也别放太重的操作。

3.3 并行汇聚:thenCombine 与 thenAcceptBoth

两个独立任务要同时执行,最后把两个结果合并成第三个结果,这种场景用thenCombine。

CompletableFuture<Integer> cartCountFuture = CompletableFuture.supplyAsync(() -> cartService.count(uid), executor); CompletableFuture<Integer> couponCountFuture = CompletableFuture.supplyAsync(() -> couponService.count(uid), executor); CompletableFuture<TotalCount> totalFuture = cartCountFuture.thenCombine(couponCountFuture, (cartCount, couponCount) -> new TotalCount(cartCount, couponCount));

它的执行时机是两个任务都完成之后,再执行合并函数。如果你只关心两个任务都完成,不需要合并结果,可以用thenAcceptBoth。

我在真实项目里其实较少用thenCombine,因为“两个任务合并”对应场景有限,更多是“一堆任务并行最后汇总”,那就要靠下面这个allOf。

3.4 多任务汇聚:allOf 与 anyOf

allOf放一串 CompletableFuture 进去,等它们全部完成,返回CompletableFuture<Void>。它本身不给结果,所以汇合时要自己遍历取值:

List<CompletableFuture<Integer>> futures = userIds.stream() .map(id -> CompletableFuture.supplyAsync(() -> userService.queryUnread(id), executor)) .collect(Collectors.toList()); CompletableFuture<List<Integer>> resultFuture = CompletableFuture .allOf(futures.toArray(new CompletableFuture[0])) .thenApply(v -> futures.stream() .map(CompletableFuture::join) .collect(Collectors.toList()));

这段代码是聚合场景的地基。注意最后用的是join()而不是get(),区别是先讲清楚:join()抛的是CompletionException,get()抛的是受检查的ExecutionException和InterruptedException。在 lambda 里用join()不用 catch 受检异常,写起来更顺手,代价是异常类型要被统一处理。

anyOf用的场景相对少,但它很擅长“去多个数据源查库存,谁先返回用谁”这类快速响应场景:

CompletableFuture<Object> anyFuture = CompletableFuture.anyOf(redisStockFuture, dbStockFuture, remoteStockFuture);

注意一点:anyOf只关心“谁先返回”,对先返回的那个值做最早响应;其余任务不会被取消,还在后台跑,这个务必要清楚,别因为用了anyOf就以为没跑完的任务不占资源。

4. 实战:一个高并发聚合接口从 0 到 1

4.1 业务场景模型:IM 未读数聚合

我把这个场景落在热词里有人搜过的“高并发 IM”上,因为它是再典型不过的异步编排场景。

假设一个 IM 应用首页要展示四类数据:

  • 会话列表未读总数,来自 Redis
  • 好友申请未读数,来自用户服务
  • 群消息未读数,来自群服务
  • 关注动态数,来自 Feed 服务

这四个数据源互相完全独立,单个平均耗时约 30ms,最慢的群服务可能要 60ms。如果串行调用,最理想也得 150ms,实际可能要 200ms 以上。首页刷新频率高,这个接口就是线上最大的热点之一。

编排思路非常清晰:四个任务全扇出,并行去查;每个任务配exceptionally兜底;最终allOf扇入汇总。整个接口目标是把 P99 控制在 80ms 以内。

4.2 完整代码与关键步骤逐段解释

先看 SpringBoot 里的 Controller 怎么写:

@RestController public class ImBadgeController { private final ThreadPoolTaskExecutor bizExecutor; public ImBadgeController(ThreadPoolTaskExecutor bizExecutor) { this.bizExecutor = bizExecutor; } @GetMapping("/im/badge") public ImBadgeVO getBadge(@RequestParam Long uid) { // 1. 并行发起四个数据源查询 CompletableFuture<Long> sessionUnreadFuture = CompletableFuture .supplyAsync(() -> unreadService.sessionUnread(uid), bizExecutor) .exceptionally(ex -> fallback("sessionUnread", ex)); CompletableFuture<Integer> friendReqFuture = CompletableFuture .supplyAsync(() -> friendService.friendReqUnread(uid), bizExecutor) .exceptionally(ex -> fallback("friendReqUnread", ex)); CompletableFuture<Integer> groupMsgFuture = CompletableFuture .supplyAsync(() -> groupService.groupMsgUnread(uid), bizExecutor) .exceptionally(ex -> fallback("groupMsgUnread", ex)); CompletableFuture<Integer> feedFuture = CompletableFuture .supplyAsync(() -> feedService.feedCount(uid), bizExecutor) .exceptionally(ex -> fallback("feedCount", ex)); // 2. 等待所有任务完成 CompletableFuture.allOf( sessionUnreadFuture, friendReqFuture, groupMsgFuture, feedFuture ).join(); // 3. 汇合结果,确保每个 Future 都有值 return ImBadgeVO.builder() .sessionUnread(safeGet(sessionUnreadFuture)) .friendReq(safeGet(friendReqFuture)) .groupMsg(safeGet(groupMsgFuture)) .feed(safeGet(feedFuture)) .build(); } private long fallback(String source, Throwable ex) { log.warn("[im-badge] {} 查询失败,使用兜底值 0", source, ex); return 0L; } private long safeGet(CompletableFuture<? extends Number> future) { try { return future.get(500, TimeUnit.MILLISECONDS).longValue(); } catch (Exception e) { log.error("[im-badge] 等待任务结果超时", e); return 0L; } } }

这段代码有几个重要细节:

第一步,四个任务都是通过同一个bizExecutor提交的,这就让它们的线程来自同一个池,后续可以通过线程池监控统一观察。

第二步,exceptionally是逐个挂上去的,不是挂在allOf之后。这是有讲究的:它能把“单个数据源失败”挡在任务内部,避免一个服务挂掉拖垮整个聚合接口。如果挂在整个链的末尾,某个异常会把整个流程都带进异常分支。

第三步,allOf().join()无参版存在隐患——所有任务都完成才返回,如果有任务没配超时,可能会等很久。所以我习惯在每个 Future 取值时用带超时的get(500, TimeUnit.MILLISECONDS)再加一层保险,这就是我写safeGet而不是直接join()的原因。

4.3 超时控制与降级策略

上面代码里future.get(500, TimeUnit.MILLISECONDS)是 JDK8 里最常见的超时手段。如果你用的是 JDK9+,还有更优雅的方式:直接在 CompletableFuture 上设置超时,超过时间自动完成并返回默认值:

CompletableFuture<Integer> feedFuture = CompletableFuture .supplyAsync(() -> feedService.feedCount(uid), bizExecutor) .completeOnTimeout(0, 500, TimeUnit.MILLISECONDS) .exceptionally(ex -> fallback("feedCount", ex));

completeOnTimeout的好处是超时后这个 CompletableFuture 会以默认值完成,不会抛异常。如果你希望超时后直接让任务失败,那用orTimeout更合适。两个方法都是 JDK9 引入的,老项目升级时要留意。

降级的核心思想特别简单:数据源偶尔失败是常态,聚合接口对单个数据源的敏感性要降到最低,把每个子任务都变成“有兜底”的状态。这也是我在生产环境上线聚合接口前一定会做的一步——把服务降级演练直接做成常规压测项目。

5. 高频翻车现场与排查避坑实录

5.1 线程池自己等自己,接口全部假死

这是我线上遇到过最诡异的问题。现象是:某个接口压测到一定程度,所有请求全部卡住,线程 dump 一看,业务线程池里的线程全在CompletableFuture.join()等子任务,子任务也提交到了同一个线程池,但线程池线程全被“等待子任务”占满,子任务永远排不上号,形成死锁。

最后排查出来的核心问题一句话就能说清:在一个线程池线程内,不要提交新的任务到同一个线程池,然后又去 join 等待它完成。你等的那个子任务可能要排队,而排队的位置已经被“正在等的你”占住了。当池里只有一个线程时,必死无疑;池里有多个线程时,也可能因为排队任务过多而大面积饿死。

规避方式有两种:一是把“编排等待”放到 Controller 请求线程里做,所有子任务都在池里并行执行,在主线程里allOf().join();二是拆两个池,编排池和业务执行池分离。实战里我更推荐前者,职责更清晰。

5.2 ThreadLocal 上下文丢失,排查半天找不到原因

异步任务一旦切到别的线程,主线程的 ThreadLocal 就带不过去了。这个坑在高并发平台项目里特别致命——你前面的过滤器刚把用户 ID、traceId 放进 ThreadLocal,异步任务里 log 就全打不出 requestId,排查问题像大海捞针。

SpringBoot 的ThreadPoolTaskExecutor提供了TaskDecorator,可以在提交任务时把上下文塞回去。这是我项目里的标准配置:

ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setTaskDecorator(runnable -> { Map<String, String> context = TraceHolder.get(); return () -> { try { TraceHolder.set(context); runnable.run(); } finally { TraceHolder.clear(); } }; });

这还不够,Spring 的@Transactional在异步任务里默认也失效,因为它同样依赖 ThreadLocal 里的事务上下文。如果你要在异步任务里做数据库写操作,别指望注解,请用编程式事务平台,或者在提交任务前把事务准备好。这是很多人开发时意识不到、上线才炸的经典问题。

5.3 异常被吞,代码里却以为处理了

CompletableFuture 的异常捕获有个反直觉的地方:链式调用里某个环节抛异常,如果后面没有exceptionally,异常会一直藏在 CompletableFuture 里,不打印、不抛给调用方。你只在主线程future.join()才知道出错了,但如果主线程也没处理,线上就是静默失败。

我见过最离谱的案例是:异步任务里分了两个分支,一个分支成功,另一个分支抛了 NPE,日志里完全没有,好几天后用户反馈某块数据总是不对,最后翻代码才发现异常被吞了。这里我给三条经验:

  • 每个异步任务的第一层就配exceptionally,哪怕只是记录日志,也要让异常有出口。
  • 日志里把业务标识打全——uid、订单号、traceId,不然线上定位无从下手。
  • 不要在whenComplete里直接返回兜底值,whenComplete不改链路上的结果,很多人误以为里面做了兜底,实际只是“看了一眼中途的值”。

5.4 从哪看线程池的健康状态

没有监控,线程池配置就是盲人摸象。好在 SpringBoot 对线程池的监控非常方便,不用额外引入复杂框架,直接看核心指标:

指标方法看什么
活跃线程数getActiveCount()长期满负载说明线程数可能偏低
当前池大小getPoolSize()观察线程是否正常扩容
队列积压getQueue().size()持续增长说明消费速度跟不上
已完成任务数getCompletedTaskCount()对比增量判断吞吐变化
拒绝任务数自定义 Count 变量这个最重要,拒绝就是丢业务

我给线程池写了个很简单的监控组件:每 30 秒把上述指标打一次 WARN 日志,队列积压超过阈值的加一条告警。这套东西不求花哨,但排障效率高到爆,一看日志就知道是不是线程池扛不住了,而不是去猜上下游。

写在最后的一点经验

这套 SpringBoot + CompletableFuture + 线程池的组合,我在生产环境维护了快两年,最大的体会是:异步编排带来的性能提升不是玄学,而是可以量化的收益;但代价是你必须对线程池和异常传播机制有足够的敬畏心。每次接入新场景,我都会先画清楚编排关系,再问自己三个问题:子任务的线程池隔离好了吗?每个子任务的超时和兜底配了吗?异常有出口、上下文能传递吗?

分享一个后续值得关注的方向:如果项目已经升级到 JDK 21,可以试试虚拟线程,它能大幅降低“线程阻塞等待 IO”的成本,和 CompletableFuture 的编排思想并不冲突。但不管底层细节怎么变,你在这套方案里学到的并发拆分思维、线程池容量规划方法、异常兜底策略,在高并发服务里会一直有用。

返回列表