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