CompletableFuture:异步编排使接口摆脱串行等待
CompletableFuture:异步编排使接口摆脱串行等待
目录
- 同步调用的痛点
- Future 的不足
- 创建异步任务
- 链式处理
- 组合多个任务
- 异常处理
- 线程池配置
- Spring Boot 实战:查商品详情
- CompletableFuture vs @Async
- 小结
同步调用的痛点
一个商品详情页的接口,需要调三个远程服务:查库存、查价格、查评论。
StockVOstock=stockService.getStock(productId);// T1PriceVOprice=priceService.getPrice(productId);// T2CommentVOcomment=commentService.getComment(productId);// T3串行执行,总耗时 = T1 + T2 + T3。但这三个调用之间没有依赖关系,完全可以同时发出去,等全部返回再组装。并行执行,总耗时 = max(T1, T2, T3) + 调度开销。
假设三个服务各 300ms,串行 900ms,并行大约 310ms 左右。它将串行变成了并行,避免了等待时间的重叠。
CompletableFuture 不会让单个任务本身变快,它优化的是等待时间。CPU 密集型任务用它就没有显著效果了,线程切换本身也有开销,下游服务也可能成为瓶颈。它解决的是这类场景:几个独立的 IO 等待,没必要排队。
Future 的不足
Java 5 就有了Future,配合线程池可以提交异步任务:
ExecutorServicepool=Executors.newFixedThreadPool(3);Future<StockVO>stockFuture=pool.submit(()->stockService.getStock(productId));Future<PriceVO>priceFuture=pool.submit(()->priceService.getPrice(productId));Future<CommentVO>commentFuture=pool.submit(()->commentService.getComment(productId));// 等结果StockVOstock=stockFuture.get();PriceVOprice=priceFuture.get();CommentVOcomment=commentFuture.get();三个任务确实是并行提交的,但Future本质上只是一个结果占位符。拿到结果后,你需要手动管理后续流程:手动等待、手动组合、手动处理异常。当业务流程简单时还能凑合,一旦任务之间有依赖关系(A 完成后再执行 B,B 和 C 的结果合并后执行 D),代码就会迅速失控。
Future的几个硬伤:
| 问题 | 说明 |
|---|---|
get()阻塞 | 必须等结果回来,当前线程被卡住 |
| 没法组合 | A 完成后再执行 B,这种链式编排做不了 |
| 没法回调 | 任务完成后想通知一下,没有回调机制 |
| 异常处理弱 | 只能try-catch get()抛出的异常 |
Java 8 引入了CompletableFuture,把这些痛点全解决了。
创建异步任务
CompletableFuture提供了两个静态方法来创建异步任务:
// 有返回值CompletableFuture<String>future=CompletableFuture.supplyAsync(()->{returnqueryFromDB();});// 没有返回值CompletableFuture<Void>future=CompletableFuture.runAsync(()->{saveToDB(data);});supplyAsync用于有返回值的场景,runAsync用于只执行动作不需要返回值的场景。
默认情况下,这两个方法用的是ForkJoinPool.commonPool()。这个公共线程池是整个 JVM 共享的,线程数 = CPU 核心数 - 1。如果某个任务阻塞了(比如调远程服务),会占住 commonPool 的线程不放,影响其他用 commonPool 的任务。生产环境建议指定自己的线程池:
ExecutorServicepool=Executors.newFixedThreadPool(10);CompletableFuture<StockVO>future=CompletableFuture.supplyAsync(()->{returnstockService.getStock(productId);},pool);笔者主页的文章:《线程池参数调优:corePoolSize 怎么设置》中有详细讲到怎么配置适合自己项目的线程池参数。
链式处理
拿到异步结果后,通常需要做进一步处理。CompletableFuture提供了一组then方法,支持链式调用。
thenApply:转换结果。类似 Stream 的map,把 A 变成 B。
CompletableFuture<String>future=CompletableFuture.supplyAsync(()->queryFromDB())// 返回原始数据.thenApply(data->format(data));// 转换成格式化后的字符串thenAccept:消费结果。拿到结果做点事情,但不返回新值。
CompletableFuture<Void>future=CompletableFuture.supplyAsync(()->queryFromDB()).thenAccept(data->log.info("查询结果: {}",data));// 打日志,不返回thenRun:执行后续动作。不关心前一步的返回值,只想在它完成后做点事。
CompletableFuture<Void>future=CompletableFuture.supplyAsync(()->queryFromDB()).thenRun(()->log.info("查询完成"));// 不关心结果,只关心"完成了"三个方法的区别:
| 方法 | 是否接收前一步结果 | 是否有返回值 | 典型用途 |
|---|---|---|---|
| thenApply | 是 | 是 | 转换数据格式 |
| thenAccept | 是 | 否 | 打日志、写缓存 |
| thenRun | 否 | 否 | 触发后续动作 |
这三个方法都有一个Async后缀的版本(thenApplyAsync、thenAcceptAsync、thenRunAsync)。不带Async的在前一个任务的线程里执行,带Async的会提交到线程池执行。带Async的版本默认也用ForkJoinPool.commonPool(),可以传第二个参数指定线程池:
// 使用自己的线程池执行异步转换future.thenApplyAsync(data->format(data),myPool);大多数场景用不带Async的就够了。只有当前一步的计算很轻、后一步很重(比如要调远程服务)时,才需要用Async版本把后续步骤丢到业务线程池里。
组合多个任务
链式处理解决的是"一个接一个"的问题。但更常见的场景是"几个任务同时跑,跑完了汇总结果"。
thenCompose:串行依赖
有时候异步任务之间有依赖:先查用户 ID,再用 ID 查订单。
CompletableFuture<List<Order>>future=getUserInfo(userId).thenCompose(user->getOrders(user.getId()));thenCompose类似 flatMap,前一步的结果作为下一步的输入,返回的还是CompletableFuture,不会嵌套成CompletableFuture<CompletableFuture<...>>。
thenCombine:合并两个独立任务
两个任务之间没有依赖,但需要把两个结果合在一起用。
CompletableFuture<User>userFuture=CompletableFuture.supplyAsync(()->queryUser(),pool);CompletableFuture<Member>memberFuture=CompletableFuture.supplyAsync(()->queryMember(),pool);// 两个任务并行执行,都完成后合并结果CompletableFuture<UserVO>result=userFuture.thenCombine(memberFuture,(user,member)->newUserVO(user,member));thenCombine和thenCompose的区别:
| 方法 | 任务关系 | 输入 | 输出 |
|---|---|---|---|
| thenCompose | 串行依赖 | 前一步的结果决定下一步做什么 | CompletableFuture |
| thenCombine | 并行独立 | 两个任务各自的结果 | CompletableFuture |
典型场景:查用户信息 + 查会员等级,合并成展示数据;查商品 + 查库存,合并成商品卡片。两个请求并行发出,都回来后再组装。
allOf:等全部完成
CompletableFuture<StockVO>stockFuture=CompletableFuture.supplyAsync(()->stockService.getStock(productId),pool);CompletableFuture<PriceVO>priceFuture=CompletableFuture.supplyAsync(()->priceService.getPrice(productId),pool);CompletableFuture<CommentVO>commentFuture=CompletableFuture.supplyAsync(()->commentService.getComment(productId),pool);// 等三个任务全部完成CompletableFuture.allOf(stockFuture,priceFuture,commentFuture).join();allOf本身不返回结果,它只负责等。要拿结果,还得从各个 Future 里取:
StockVOstock=stockFuture.join();PriceVOprice=priceFuture.join();CommentVOcomment=commentFuture.join();这里用join()而不是get(),因为join()不抛受检异常,代码更简洁。实际效果一样,都是阻塞等待。
anyOf:谁先完成用谁
CompletableFuture<String>f1=CompletableFuture.supplyAsync(()->queryFromRedis());// Redis 快,可能先返回CompletableFuture<String>f2=CompletableFuture.supplyAsync(()->queryFromDB());// DB 慢,可能后返回Objectresult=CompletableFuture.anyOf(f1,f2).join();anyOf返回第一个完成的任务的结果。适合"多级缓存"的场景:先查 Redis,同时查 DB,谁先返回用谁。
异常处理
异步任务抛了异常,get()或join()时会抛CompletionException。如果不处理,异常就静默丢了,排查问题时根本不知道哪里出了错。
exceptionally:异常时的兜底
CompletableFuture<String>future=CompletableFuture.supplyAsync(()->{if(somethingWrong)thrownewRuntimeException("出错了");return"正常结果";}).exceptionally(ex->{log.error("异步任务失败: {}",ex.getMessage());return"默认值";// 返回一个兜底值});exceptionally相当于 catch,给一个兜底的返回值。
handle:统一处理正常和异常
CompletableFuture<String>future=CompletableFuture.supplyAsync(()->queryFromDB()).handle((result,ex)->{if(ex!=null){log.error("查询失败: {}",ex.getMessage());return"默认值";}returnresult;});handle同时接收正常结果和异常对象,比exceptionally更灵活。正常时ex为 null,异常时result为 null。
子任务异常隔离
用allOf编排多个任务时,一个容易踩的坑:如果某个子任务抛了异常,allOf也会异常,后面再调join()还是会抛。allOf的exceptionally只处理allOf本身的异常,不处理子任务的异常。
正确的做法是每个子任务自己兜底,allOf只负责编排:
CompletableFuture<StockVO>stockFuture=CompletableFuture.supplyAsync(()->stockService.getStock(productId),pool).exceptionally(ex->{log.error("库存查询失败",ex);returndefaultStock;// 返回默认库存});CompletableFuture<PriceVO>priceFuture=CompletableFuture.supplyAsync(()->priceService.getPrice(productId),pool).exceptionally(ex->{log.error("价格查询失败",ex);returndefaultPrice;});CompletableFuture<CommentVO>commentFuture=CompletableFuture.supplyAsync(()->commentService.getComment(productId),pool).exceptionally(ex->{log.error("评论查询失败",ex);returndefaultComment;});// allOf 只管编排,不管异常(子任务已经各自兜底了)CompletableFuture.allOf(stockFuture,priceFuture,commentFuture).join();ProductDetailVOdetail=newProductDetailVO();detail.setStock(stockFuture.join());// 拿到的是正常结果或兜底值detail.setPrice(priceFuture.join());detail.setComment(commentFuture.join());这样即使价格服务挂了,库存和评论正常返回,接口照样能用,只是价格显示默认值。
线程池配置
前面提到过,生产环境不要直接用Executors.newFixedThreadPool()。原因是它内部用的是无界队列LinkedBlockingQueue:
// 不推荐:队列无界,高流量下任务堆积导致 OOMExecutorServicepool=Executors.newFixedThreadPool(10);Spring Boot 项目推荐用ThreadPoolTaskExecutor,可以控制核心线程数、最大线程数和队列容量:
@BeanpublicThreadPoolTaskExecutorasyncExecutor(){ThreadPoolTaskExecutorexecutor=newThreadPoolTaskExecutor();executor.setCorePoolSize(10);// 核心线程数executor.setMaxPoolSize(20);// 最大线程数executor.setQueueCapacity(100);// 队列容量,超出后创建新线程executor.setThreadNamePrefix("async-");executor.setRejectedExecutionHandler(newThreadPoolExecutor.CallerRunsPolicy());executor.initialize();returnexecutor;}然后注入使用:
@Autowired@Qualifier("asyncExecutor")privateThreadPoolTaskExecutorasyncExecutor;CompletableFuture<StockVO>future=CompletableFuture.supplyAsync(()->stockService.getStock(productId),asyncExecutor);这样当队列满时会触发拒绝策略(这里用的是CallerRunsPolicy,提交者自己执行),不会无限堆积。线程池的行为是可预期、可监控的。
Spring Boot 实战:查商品详情
把前面的知识串起来,写一个完整的例子:并行查库存、价格、评论,每个子任务独立兜底,合并成商品详情返回。
@ServicepublicclassProductService{@AutowiredprivateStockServicestockService;@AutowiredprivatePriceServicepriceService;@AutowiredprivateCommentServicecommentService;@Autowired@Qualifier("asyncExecutor")privateThreadPoolTaskExecutorasyncExecutor;publicProductDetailVOgetProductDetail(LongproductId){// 三个任务并行提交,每个任务独立处理异常CompletableFuture<StockVO>stockFuture=CompletableFuture.supplyAsync(()->stockService.getStock(productId),asyncExecutor).exceptionally(ex->{log.error("库存查询失败: {}",ex.getMessage());returnStockVO.empty();});CompletableFuture<PriceVO>priceFuture=CompletableFuture.supplyAsync(()->priceService.getPrice(productId),asyncExecutor).exceptionally(ex->{log.error("价格查询失败: {}",ex.getMessage());returnPriceVO.empty();});CompletableFuture<CommentVO>commentFuture=CompletableFuture.supplyAsync(()->commentService.getComment(productId),asyncExecutor).exceptionally(ex->{log.error("评论查询失败: {}",ex.getMessage());returnCommentVO.empty();});// 等全部完成CompletableFuture.allOf(stockFuture,priceFuture,commentFuture).join();// 组装结果(此时不会阻塞,因为 allOf 已经等完了)ProductDetailVOdetail=newProductDetailVO();detail.setProductId(productId);detail.setStock(stockFuture.join());detail.setPrice(priceFuture.join());detail.setComment(commentFuture.join());returndetail;}}三个远程调用并行发出,总耗时取决于最慢的那个。某个服务挂了不影响整体,接口降级返回默认值。
再看一个稍微进阶的场景:先查用户信息,再并行查该用户的订单和优惠券。
publicUserDashboardVOgetDashboard(LonguserId){// 第一步:查用户信息(后续步骤依赖它)CompletableFuture<UserVO>userFuture=CompletableFuture.supplyAsync(()->userService.getUser(userId),asyncExecutor);// 第二步:拿到用户信息后,并行查订单和优惠券CompletableFuture<List<Order>>ordersFuture=userFuture.thenComposeAsync(user->CompletableFuture.supplyAsync(()->orderService.getOrders(user.getId()),asyncExecutor),asyncExecutor);CompletableFuture<List<Coupon>>couponsFuture=userFuture.thenComposeAsync(user->CompletableFuture.supplyAsync(()->couponService.getCoupons(user.getId()),asyncExecutor),asyncExecutor);// 等全部完成CompletableFuture.allOf(ordersFuture,couponsFuture).join();// 组装UserDashboardVOdashboard=newUserDashboardVO();dashboard.setUser(userFuture.join());dashboard.setOrders(ordersFuture.join());dashboard.setCoupons(couponsFuture.join());returndashboard;}第一步查用户是串行的(后面依赖它的结果),第二步查订单和优惠券是并行的。整个过程的总耗时 = 查用户的时间 + max(查订单, 查优惠券)。
CompletableFuture vs @Async
Spring 开发者可能会问:Spring 不是有@Async注解吗,为什么还要用 CompletableFuture?
// @Async 的用法@AsyncpublicFuture<StockVO>getStockAsync(LongproductId){StockVOstock=stockService.getStock(productId);returnnewAsyncResult<>(stock);}@Async解决的是"这个方法异步执行",CompletableFuture解决的是"多个异步任务如何组织"。
| 维度 | CompletableFuture | @Async |
|---|---|---|
| 定位 | 异步流程编排 | 简单异步执行 |
| 链式组合 | 强,支持 thenApply/thenCompose | 弱,需要手动管理 |
| 多个任务组合 | allOf / anyOf / thenCombine | 需要自己写等待逻辑 |
| 异常处理 | exceptionally / handle | 依赖 AsyncUncaughtExceptionHandler |
如果你只是想让某个方法异步执行,不关心结果编排,@Async够用。但如果需要多个异步任务并行、串行、合并、竞争,CompletableFuture是更合适的选择。两者不冲突,很多项目里会同时用:@Async负责把方法变成异步,返回CompletableFuture,再用 CompletableFuture 的能力做编排。
小结
CompletableFuture 解决的是"多个异步任务怎么编排"的问题。链式调用处理串行依赖,thenCombine合并独立任务的结果,allOf处理并行汇总,anyOf处理竞争场景,handle和exceptionally兜底异常。很多接口性能问题,并不是代码执行慢,而是在等待多个没有依赖关系的操作串行完成。把这些等待时间重叠起来,往往比单纯优化某一段代码更有效。