尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

CompletableFuture顺序工作流异步执行与异常处理实践

CompletableFuture顺序工作流异步执行与异常处理实践 在开发中我们经常会遇到一类需求多个环节必须按照固定顺序执行前一步成功后才能继续下一步但每一步又比较耗时比如发短信、调外部接口、写日志。如果直接在请求线程里一步步同步等待一个接口的响应时间就会变成所有步骤耗时的总和并发一高线程直接被拖垮。如果贸然改成异步又面临两个新问题怎么保证“顺序”某个环节出错之后后面还该不该继续执行这篇文章要讲的就是顺序工作流异步执行。我会用 Java 的CompletableFuture作为主线说明如何把一条有序的任务链跑在异步线程池里既保证前后依赖关系又能优雅地把控异常传播。特别是很多人踩过的坑当某个异步任务抛出异常后后续任务默认不会继续执行这时候到底应该中断还是降级还是跳过必须有一个明确的策略。如果你正在写订单流程、审批流、数据同步链路或者任何“有前后依赖的多步操作”这篇文章值得读完并收藏。我会从概念、代码、异常处理到生产实践完整拆解这一套写法。1. 顺序工作流异步执行到底解决什么问题先看一个真实场景。假设用户下单后我们需要做四件事校验商品库存锁定库存创建订单发送通知消息。这四个步骤存在明显的先后依赖必须先校验库存才能锁定库存必须锁定库存成功才能创建订单。如果“校验库存”还没结束就去“创建订单”结果大概率是订单创建了库存却没扣掉超卖问题就会出现。传统写法是同步执行public void createOrder(OrderRequest request) { boolean available inventoryService.checkStock(request.getSkuId()); if (!available) { throw new BusinessException(库存不足); } boolean locked inventoryService.lockStock(request.getSkuId(), request.getQuantity()); if (!locked) { throw new BusinessException(锁定库存失败); } Order order orderService.createOrder(request); messageService.sendMessage(order); }这种代码逻辑没问题而且容易理解。但问题在于如果每一步都耗时 200ms整个接口就需要 800ms 才能返回期间请求线程一直被占用。在 Tomcat 默认 200 线程的场景下每秒最多只能处理 250 个这样的请求而且大部分时间线程都阻塞在等待 IO 上。异步执行的价值在于让请求线程快速返回真正耗时的步骤交给后台线程池去处理。但“异步”和“顺序”看起来是矛盾的——异步是并行发出去的顺序要求一个一个来。CompletableFuture正好提供了这样的能力它允许你把有依赖的步骤串成一条执行链前一个步骤的完成结果可以直接作为下一步的输入天然保证顺序同时每个步骤都运行在线程池中。所以顺序工作流异步执行解决的本质问题不是“谁快谁慢”的微观性能而是在保持业务依赖顺序不变的前提下把同步阻塞模型替换成异步回调模型从而降低请求线程占用时间提高系统吞吐量。这篇文章适合以下读者需要优化接口响应时间但任务之间存在强依赖正在学习或使用CompletableFuture希望搞清楚它如何编排异步任务在项目里遇到“异步任务异常后不执行后续任务”的困惑想了解如何把一条业务链路拆解成可监控、可降级的异步工作流。2. 基础概念顺序工作流、异步执行与 CompletableFuture2.1 什么是顺序工作流顺序工作流指的是一组任务按照固定的先后次序执行前一个任务的输出是后一个任务的输入或者至少前一个任务的成功是后一个任务的启动条件。它的特征有三个依赖关系任务 B 必须在任务 A 完成后启动数据传递A 的结果可能传递给 B失败传播如果 A 失败B 通常不应该继续执行除非业务上允许降级。在 Java 8 之前想让多线程按顺序执行通常用ExecutorService配合Future轮询或者用CountDownLatch等待多个任务完成。这些方式要么写起来啰嗦要么很难表达“前一个结果传给后一个”的语义。2.2 什么是异步执行异步执行是一种编程模型调用方发起一个任务后不立即等待结果而是由另一个线程去执行调用方通过回调、轮询或阻塞获取最终结果。异步执行的核心收益是减少线程空闲等待时间。对于 IO 密集型操作真正消耗 CPU 的时间极少大部分时间都在等待网络、数据库、第三方接口响应。如果每个操作都占用一个线程同步等待线程资源就会被浪费。但异步执行也带来了复杂度缺少调用栈上下文排错困难异常处理方式从 try-catch 变成回调顺序依赖难以直观表达容易引入并发问题。CompletableFuture解决的就是后两个问题。2.3 CompletableFuture 的核心能力CompletableFuture是java.util.concurrent包下的一个类在 Java 8 中引入。它实现了Future和CompletionStage接口本质上是一个“可手动完成的 Future”同时支持把多个阶段串联成执行链。要用好CompletableFuture实现顺序工作流必须分清几个核心方法方法类型作用是否支持串行依赖异常处理supplyAsync静态方法提交一个带返回值的异步任务作为起点异常会记录到返回的 Future 中thenApply实例方法上一个任务完成后把结果作为参数执行新任务返回新结果是异常会继续向下传播thenAccept实例方法上一个任务完成后消费结果无返回值是异常继续传播thenCompose实例方法上一个任务完成后返回一个新的 CompletionStage用于扁平化异步链是异常继续传播exceptionally实例方法只有在上游出现异常时触发返回降级结果会截断异常传播处理掉异常handle实例方法无论成功失败都会触发接收结果和异常两个参数会截断异常传播可手动处理whenComplete实例方法无论成功失败都会触发但会保留异常继续向下传播不截断异常只能感知不能处理这里最关键的一点是thenApply、thenAccept、thenCompose这些方法默认具有“异常短路”特性。也就是说如果链上的某一个任务抛出了异常那么后续的thenApply等任务将不会执行异常会沿着链向下传播直到被exceptionally或handle处理。这就是热搜词里说的“CompletableFuture 异常后不在执行其他的异步任务”。很多新手会在这个地方卡住明明在thenApply里捕获了异常为什么下一步还是没执行原因是你没有用对处理异常的阶段方法。3. 从同步到异步先构建最小可运行环境在动手写代码之前先明确环境要求。3.1 环境说明JDK 版本Java 8 及以上本文示例使用 Java 11 验证构建工具Maven 3.6也可以用 Gradle但本文以 Maven 为例无需引入任何第三方依赖CompletableFuture是 JDK 自带能力建议安装 IntelliJ IDEA 或 Eclipse方便断点调试和观察线程变化。实际上CompletableFuture不依赖 Spring任何 Java 项目都可以直接使用。不过在实际业务项目中通常会结合 Spring 的Async或自定义线程池一起使用这部分我会在后面的最佳实践里说明。3.2 创建一个简单的 Maven 项目如果是新建项目pom.xml只需要最基础的配置?xml version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion groupIdcom.example/groupId artifactIdasync-workflow-demo/artifactId version1.0-SNAPSHOT/version packagingjar/packaging properties maven.compiler.source11/maven.compiler.source maven.compiler.target11/maven.compiler.target project.build.sourceEncodingUTF-8/project.build.sourceEncoding /properties /project如果你使用的是已有的 Spring Boot 项目完全不需要额外添加依赖。CompletableFuture在java.base模块中直接 import 使用即可。4. 用 CompletableFuture 实现顺序工作流接下来我们用一个贴近业务的示例来演示模拟“下单成功后依次执行库存校验、库存锁定、订单创建、发送通知”四个步骤。为了模拟耗时每个步骤都会让当前线程休眠几百毫秒并打印线程名和结果。这样你就能清楚看到任务在哪个线程上执行以及顺序是否正确。4.1 定义任务服务先定义一个模拟业务服务的类每个方法都返回一个结果同时支持抛出异常。这里把异常处理逻辑放在后面演示所以第一个版本我们先保证正常流程。package com.example.asyncworkflow; import java.util.concurrent.TimeUnit; /** * 模拟业务服务 */ public class OrderWorkflowService { /** * 校验库存 */ public String checkStock(String skuId) { sleep(300); System.out.println([checkStock] 校验商品库存 skuId skuId , 线程 Thread.currentThread().getName()); return 库存充足; } /** * 锁定库存 */ public String lockStock(String skuId, int quantity) { sleep(400); System.out.println([lockStock] 锁定库存 skuId skuId , quantity quantity , 线程 Thread.currentThread().getName()); return 库存锁定成功; } /** * 创建订单 */ public String createOrder(String skuId, int quantity) { sleep(500); System.out.println([createOrder] 创建订单 skuId skuId , quantity quantity , 线程 Thread.currentThread().getName()); return ORD-20250101-001; } /** * 发送通知 */ public String sendMessage(String orderId) { sleep(200); System.out.println([sendMessage] 发送订单通知 orderId orderId , 线程 Thread.currentThread().getName()); return 通知发送成功; } private void sleep(long millis) { try { TimeUnit.MILLISECONDS.sleep(millis); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException(任务被中断, e); } } }这里简单说明一下每个方法内部的sleep模拟外部调用耗时打印线程名是为了观察任务运行在哪条线程上。注意Thread.currentThread().interrupt()是处理中断的标准做法避免吞掉中断状态。4.2 使用 thenApply 串起整条链现在我们用CompletableFuture把四个步骤串起来。package com.example.asyncworkflow; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class SequentialWorkflowDemo { public static void main(String[] args) throws Exception { // 创建一个固定大小为 4 的线程池 ExecutorService executor Executors.newFixedThreadPool(4); OrderWorkflowService service new OrderWorkflowService(); CompletableFutureString future CompletableFuture .supplyAsync(() - service.checkStock(SKU-1001), executor) .thenApply(result - { System.out.println([Step1] result); return service.lockStock(SKU-1001, 2); }) .thenApply(result - { System.out.println([Step2] result); return service.createOrder(SKU-1001, 2); }) .thenApply(result - { System.out.println([Step3] 订单号: result); return service.sendMessage(result); }); // 等待整个流程结束并获取最终结果 String finalResult future.get(); System.out.println([Final] finalResult); executor.shutdown(); } }运行这段代码输出如下线程名可能不同[Step1] 库存充足 [Step2] 库存锁定成功 [Step3] 订单号: ORD-20250101-001 [Final] 通知发送成功 [checkStock] 校验商品库存 skuIdSKU-1001, 线程pool-1-thread-1 [lockStock] 锁定库存 skuIdSKU-1001, quantity2, 线程pool-1-thread-1 [createOrder] 创建订单 skuIdSKU-1001, quantity2, 线程pool-1-thread-1 [sendMessage] 发送订单通知 orderIdORD-20250101-001, 线程pool-1-thread-1注意看四个步骤全部运行在pool-1-thread-1上。这是因为thenApply默认的异步编排策略如果前一个任务已经完成那么下一个thenApply可能会在调用线程上执行如果前一个任务未完成则会在前一个任务所在的线程上继续执行。在我们的场景中每个阶段都是连续提交的前一个阶段没有完成所以后续阶段复用了同一个线程。这一点并不需要过度担心关键是“顺序”被保证了后面的步骤一定在接收到前一个步骤的返回值之后才开始。4.3 更灵活的 thenCompose 写法thenApply的返回值会直接包装成新的CompletableFuture如果内部还要返回新的异步任务就会出现CompletableFutureCompletableFutureT的嵌套结构。为了避免这种嵌套可以使用thenCompose。CompletableFutureString future2 CompletableFuture .supplyAsync(() - service.checkStock(SKU-2001), executor) .thenCompose(checkResult - { System.out.println([Step1] checkResult); // 返回一个新的异步阶段 return CompletableFuture.supplyAsync( () - service.lockStock(SKU-2001, 3), executor); }) .thenCompose(lockResult - { System.out.println([Step2] lockResult); return CompletableFuture.supplyAsync( () - service.createOrder(SKU-2001, 3), executor); }) .thenApply(orderId - { System.out.println([Step3] 订单号: orderId); return service.sendMessage(orderId); });这段代码的运行结果与上一版一致。区别在于thenCompose要求函数返回一个CompletionStage它会把返回的 Stage 自动展开适合每个步骤都是独立异步任务的场景。实际项目中如果某个步骤需要单独控制线程池或做额外监控用thenCompose会更清晰。如果只是简单的同步计算和消费用thenApply就够了。5. 异常处理为什么后续异步任务不执行了现在进入本文最关键的部分。先看一个现象。我们把lockStock方法改成可能抛异常public String lockStock(String skuId, int quantity) { sleep(400); if (ERROR-SKU.equals(skuId)) { throw new RuntimeException(库存服务异常); } System.out.println([lockStock] 锁定库存 skuId skuId , quantity quantity , 线程 Thread.currentThread().getName()); return 库存锁定成功; }然后使用原来的链式调用传入ERROR-SKU看看会发生什么。CompletableFutureString futureError CompletableFuture .supplyAsync(() - service.checkStock(ERROR-SKU), executor) .thenApply(result - { System.out.println([Step1] result); return service.lockStock(ERROR-SKU, 2); }) .thenApply(result - { System.out.println([Step2] result); return service.createOrder(ERROR-SKU, 2); }) .thenApply(result - { System.out.println([Step3] 订单号: result); return service.sendMessage(result); });运行后你会发现控制台只打印了[Step1] 库存充足然后就没有任何后续输出了。最终调用futureError.get()时会抛出ExecutionException根本原因是RuntimeException: 库存服务异常。这就是热搜词描述的“CompletableFuture 异常后不再执行其他的异步任务”。从设计角度看这是合理的默认行为业务链上某一步失败说明后续步骤的前提不再成立继续执行只会产生脏数据。但在实际工程中我们往往需要根据业务场景决定后续步骤全部取消整个流程标记失败对异常进行降级给一个默认值让流程继续走跳过失败步骤允许非关键步骤发生异常时不影响主链路。这三种策略分别对应exceptionally、handle和“在具体步骤里自行捕获”。下面逐一说明。5.1 使用 exceptionally 处理异常并降级exceptionally只在链上出现异常时触发它接收异常对象并返回一个降级结果。这个结果会替换掉异常后续步骤会继续执行。CompletableFutureString futureWithFallback CompletableFuture .supplyAsync(() - service.checkStock(ERROR-SKU), executor) .thenApply(result - { System.out.println([Step1] result); return service.lockStock(ERROR-SKU, 2); }) .exceptionally(ex - { System.out.println([Exception] 捕获异常: ex.getMessage()); // 降级返回一个默认锁定结果 return 库存锁定失败走降级逻辑; }) .thenApply(result - { System.out.println([Step2] result); return service.createOrder(ERROR-SKU, 2); }) .thenApply(result - { System.out.println([Step3] 订单号: result); return service.sendMessage(result); });需要注意exceptionally放在哪个位置决定它能捕获到哪一段的异常。上面这个例子中exceptionally位于lockStock之后所以它只能捕获checkStock和lockStock之间的异常。如果createOrder也抛异常这个exceptionally捕获不到需要再往下游添加新的处理节点。另外exceptionally一旦返回结果异常就被“吞掉”了下游看到的是正常结果。这意味着你必须确保降级结果不会污染后续业务。比如库存锁定失败但下游仍然创建订单这是非常危险的。所以在使用exceptionally时要清楚业务边界只有允许“降级继续”的步骤才适合这样处理。5.2 使用 handle 同时处理正常结果和异常handle方法与exceptionally不同它无论上游成功还是失败都会执行。它接收两个参数正常结果和异常对象两者至少有一个为 null。通过判断异常对象是否为空可以决定是返回正常处理结果还是降级结果。CompletableFutureString futureWithHandle CompletableFuture .supplyAsync(() - service.checkStock(HANDLE-SKU), executor) .thenApply(result - { System.out.println([Step1] result); return service.lockStock(HANDLE-SKU, 2); }) .handle((lockedResult, ex) - { if (ex ! null) { System.out.println([Handle] 捕获异常: ex.getMessage()); return LOCK_FAILED; } System.out.println([Handle] 正常锁定: lockedResult); return lockedResult; }) .thenApply(result - { System.out.println([Step2] 处理结果: result); if (LOCK_FAILED.equals(result)) { // 业务上决定失败后不再创建订单 throw new IllegalStateException(前置条件未满足流程终止); } return service.createOrder(HANDLE-SKU, 2); }) .thenApply(orderId - { System.out.println([Step3] 订单号: orderId); return service.sendMessage(orderId); });这种写法的好处是你能在同一处同时处理成功和失败分支逻辑更集中。但要注意handle返回的结果仍然会继续传递到下游如果你希望“失败后中断整条链”那么在handle内部要么重新抛出异常要么返回一个下游能识别的特殊值并让下游显式抛出新异常。前者更直接handle里如果抛出异常这个异常会替换原来的异常向下传递。哪种方式更好我的判断是如果只是简单降级用exceptionally如果需要对成功和失败结果做统一转换用handle如果只是想记录日志不做任何结果替换用whenComplete。但whenComplete不会截断异常日志记录后异常还会继续向下传播这一点务必记住。5.3 只想要“感知异常”不想处理whenComplete方法在阶段完成时触发无论成功失败。它接收结果和异常但它的返回值类型是CompletableFutureT其中 T 是上游的结果类型。它的特点是如果传入参数有异常whenComplete执行完后异常仍然会继续向下传播除非你在回调里抛出别的异常。CompletableFutureString futureWithWhenComplete CompletableFuture .supplyAsync(() - service.checkStock(WC-SKU), executor) .thenApply(result - { System.out.println([Step1] result); return service.lockStock(WC-SKU, 2); }) .whenComplete((result, ex) - { if (ex ! null) { System.out.println([WhenComplete] 执行到此时出现异常: ex.getMessage()); // 这里可以记录日志、发送告警但不会吞掉异常 } else { System.out.println([WhenComplete] 当前结果: result); } }) .thenApply(result - { // 如果上游有异常这个阶段不会执行 System.out.println([Step2] 只有成功才能到这里 result); return service.createOrder(WC-SKU, 2); });运行这段代码如果lockStock抛异常whenComplete会被触发并打印日志但随后thenApply不会执行。这就引出一个非常实用的结论“异常后不再执行其他异步任务”并不是一个 bug而是 CompletableFuture 默认的短路机制。你要做的不是想方设法绕过它而是根据业务规则在合适的位置用exceptionally、handle或whenComplete接管异常流。6. 完整示例一个带异常策略的顺序工作流下面我们把上面的知识整合成一个相对完整的示例。这个示例模拟真实业务库存不足时直接失败库存服务异常时走降级重试订单创建失败时发送告警并终止通知发送失败时不影响主流程。package com.example.asyncworkflow; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; public class WorkflowWithExceptionStrategy { public static void main(String[] args) throws Exception { ExecutorService executor Executors.newFixedThreadPool(8); OrderWorkflowService service new OrderWorkflowService(); // 模拟库存服务异常一次的 SKU String skuId RETRY-SKU; CompletableFutureString finalResult CompletableFuture // 第一步校验库存同步等待结果这里必须成功 .supplyAsync(() - service.checkStock(skuId), executor) // 第二步锁定库存失败后最多重试一次 .thenCompose(checkResult - { System.out.println([Workflow] 校验结果: checkResult); return lockStockWithRetry(service, skuId, 2, executor); }) // 第三步创建订单失败则整个流程终止 .thenCompose(lockResult - { System.out.println([Workflow] 锁定结果: lockResult); CompletableFutureString orderFuture CompletableFuture.supplyAsync(() - service.createOrder(skuId, 2), executor); return orderFuture.exceptionally(ex - { System.out.println([Workflow] 创建订单失败流程终止: ex.getMessage()); throw new IllegalStateException(创建订单失败, ex); }); }) // 第四步发送通知失败不向上抛 .thenApply(orderId - { System.out.println([Workflow] 订单号: orderId); try { return service.sendMessage(orderId); } catch (Exception e) { System.out.println([Workflow] 通知发送失败忽略: e.getMessage()); return 通知发送失败(已忽略); } }); try { String result finalResult.get(5, TimeUnit.SECONDS); System.out.println([Workflow] 最终结果: result); } catch (Exception e) { System.out.println([Workflow] 流程异常: e.getCause().getMessage()); } finally { executor.shutdown(); } } /** * 锁定库存支持失败重试一次 */ private static CompletableFutureString lockStockWithRetry( OrderWorkflowService service, String skuId, int quantity, ExecutorService executor) { CompletableFutureString firstAttempt CompletableFuture.supplyAsync(() - service.lockStock(skuId, quantity), executor); return firstAttempt.exceptionally(ex - { System.out.println([Retry] 第一次锁定失败原因: ex.getMessage() 开始重试); return service.lockStock(skuId, quantity); }); } }为了演示重试效果你可以临时修改OrderWorkflowService.lockStock让它对某个 SKU 第一次抛异常第二次成功。例如private static int lockCount 0; public String lockStock(String skuId, int quantity) { sleep(400); lockCount; if (RETRY-SKU.equals(skuId) lockCount 1) { throw new RuntimeException(库存服务暂时不可用); } System.out.println([lockStock] 锁定库存 skuId skuId , quantity quantity , 线程 Thread.currentThread().getName()); return 库存锁定成功; }运行后你会看到输出顺序大致如下[Workflow] 校验结果: 库存充足 [Retry] 第一次锁定失败原因: java.lang.RuntimeException: 库存服务暂时不可用开始重试 [Workflow] 锁定结果: 库存锁定成功 [Workflow] 订单号: ORD-20250101-001 [Workflow] 最终结果: 通知发送成功这个示例展示了三个关键理念关键步骤失败时可以针对该步骤做局部重试不可降级的业务哪怕在exceptionally里也要重新抛出异常非关键的通知步骤单独 try-catch 不会影响主链路。7. 运行结果与验证方式前面所有的代码运行起来后验证的重点不是“能打印出来”而是验证几个隐藏语义顺序是否严格保证异常后后续阶段是否按预期执行或不执行降级结果是否不会污染后续业务线程池使用是否正常有没有线程泄漏。7.1 验证顺序可以在每个步骤打印时间戳对比阶段耗时long start System.currentTimeMillis(); ... System.out.println([Time] 当前耗时: (System.currentTimeMillis() - start) ms, 步骤: xxx);如果每个步骤的sleep时长不同而最终总耗时约等于四个步骤耗时之和说明任务是顺序执行的没有并发乱序。7.2 验证异常短路把lockStock改成必抛异常观察输出。预期结果thenApply链上位于异常之后的阶段不会执行future.get()抛出ExecutionException如果添加了exceptionally异常被捕获后续阶段继续执行。7.3 验证线程池任务状态在确定不会再提交新任务后调用executor.shutdown()并检查是否有线程一直不退出。如果业务中误用了没有关闭的线程池JVM 进程可能无法正常结束。executor.shutdown(); boolean terminated executor.awaitTermination(3, TimeUnit.SECONDS); System.out.println(线程池是否终止: terminated); if (!terminated) { executor.shutdownNow(); }这段代码确保线程池在合理时间内关闭防止资源泄漏。7.4 失败时第一步看哪里如果运行结果不符合预期按以下顺序排查看第一个异常出现的位置是不是上游任务本身抛了异常看异常处理方法的放置位置是否覆盖到了异常发生点看返回值异常处理后返回的降级值是否在下游被误当成成功值看线程池队列会不会因为任务积压导致某个阶段长期不执行。8. 常见问题与排查思路以下是我在项目中遇到过的真实问题和解决思路整理成表格方便收藏。问题现象可能原因排查方式解决方案异步任务抛异常后后续阶段没有任何输出thenApply/thenCompose不具备异常捕获能力异常短路向传播在调用get()的位置打印ExecutionException查看 cause在合适位置加exceptionally或handle处理异常exceptionally捕获了异常但下游仍然失败降级返回的值不符合下游业务预期导致业务层主动抛错检查降级返回值检查下游是否对该值有判断设计降级协议值下游显式判断并决定是否继续结果表现为并行执行顺序错乱多个 stage 在thenApply中又触发独立异步任务未使用thenCompose连接查看代码中是否出现CompletableFuture.supplyAsync嵌套用thenCompose扁平化异步链future.get()一直阻塞某一步递归等待自己或者线程池核心线程耗尽使用带超时的get(timeout)查看线程池活跃线程数给所有异步等待加超时避免永久阻塞主线程调用future.get()后接口响应仍然很慢本质是同步等待异步结果异步优势被抵消观察接口耗时是否约等于工作流总耗时如果接口必须返回结果考虑前置返回后再异步执行如果必须同步拿结果使用异步并不能降低总耗时异常被吞掉日志里没有任何记录exceptionally或handle中只返回值没有记录日志在异常处理方法中增加日志输出统一异常日志记录配合链路追踪 TraceId应用关闭时线程池没有退出线程池被Executors创建后未关闭检查 JVM 线程列表查看非守护线程使用 Spring 管理生命周期或在销毁钩子中优雅关闭这些问题的核心仍然是对CompletableFuture的执行链和异常传播机制理解不透。建议你在本地多写几个最小示例把exceptionally、handle、whenComplete放在链的不同位置观察结果比看十篇文章都有效。9. 生产环境最佳实践与工程建议CompletableFuture本身并不复杂复杂的是把它放进真正的业务系统里。下面是几点工程建议。9.1 不要用默认的 ForkJoinPool 执行异步任务CompletableFuture.supplyAsync如果不传线程池默认使用ForkJoinPool.commonPool()。这个公共池的并行度默认和 CPU 核数相关容易被其他框架的异步任务抢占也容易因为阻塞操作耗尽线程。生产环境一定要传入自定义线程池。建议使用ThreadPoolExecutor手动创建并设置合理的参数ExecutorService workflowExecutor new ThreadPoolExecutor( 8, 16, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue(1000), new ThreadFactory() { Override public Thread newThread(Runnable r) { Thread t new Thread(r); t.setName(workflow-executor- t.getId()); t.setDaemon(true); return t; } }, new ThreadPoolExecutor.CallerRunsPolicy() );参数说明核心线程数按该工作流瞬间积压的任务量估算队列使用有界队列避免内存溢出拒绝策略CallerRunsPolicy在任务过多时由提交线程执行这种降级可以避免任务丢失但要注意主线程会被阻塞线程命名强烈建议否则排查问题时看不到是哪个业务线程。9.2 区分“关键步骤”和“非关键步骤”在顺序工作流里并不是每一步都必须成功。我的建议是提前画一张依赖表步骤是否关键失败策略校验库存关键直接失败后续不执行锁定库存关键重试一次仍失败则终止创建订单关键失败发送告警流程终止发送通知非关键记录日志不影响主流程有了这张表代码的异常处理位置就非常清楚了关键步骤用exceptionally包装成业务失败非关键步骤在方法内部 try-catch不向上传递。9.3 统一超时控制CompletableFuture本身没有提供orTimeout方法直到 Java 9 才加入。如果你使用 Java 8可以借助get(timeout, unit)来兜底但这会阻塞当前线程。更好的方式是用ScheduledExecutorService加completeExceptionally实现超时中断但这种写法略复杂。如果项目在 Java 9直接使用orTimeoutCompletableFutureString futureWithTimeout future .orTimeout(3, TimeUnit.SECONDS) .exceptionally(ex - { System.out.println([Timeout] 任务超时或异常: ex.getMessage()); return TIMEOUT_FALLBACK; });如果项目仍在 Java 8我的建议是尽量在get时加超时同时保证整个异步链不会无限等待。9.4 链路追踪和日志关联异步链的一个麻烦是日志打印在多条线程上排查时很难把同一笔订单的日志串起来。解决方法是给每一步传递同一个 TraceId。String traceId UUID.randomUUID().toString(); CompletableFuture .supplyAsync(() - { MDC.put(traceId, traceId); return service.checkStock(skuId); }, executor) .thenApply(result - { // 此时 MDC 不一定能继承因为线程可能变化 return service.lockStock(skuId, 2); });注意CompletableFuture不会自动传递 MDC 上下文每次进入新的thenApply阶段都可能切换到另一个线程。如果链路追踪工具没有提供线程池包装就需要手动传递。常见的做法是使用transmittable-thread-local这类工具或者每进入一个阶段就重新设置 TraceId。9.5 监控与告警每一步的耗时、成功失败次数、重试次数都应该记录下来。生产环境推荐至少记录每个阶段的耗时分布每个阶段失败次数重试次数及重试结果线程池队列长度、活跃线程数、拒绝任务数。这些数据可以导出到 Prometheus 等监控系统。如果发现某个阶段耗时持续增长就需要考虑是不是外部依赖变慢而不是盲目加线程池大小。9.6 谨慎使用 thenApply 中的耗时操作thenApply的任务如果本身是耗时操作最好单独声明为异步任务并通过thenCompose连接。否则虽然你用了异步线程池但整个链还是在一个线程上串行执行如果该线程被某个慢操作卡住其他任务也会受影响。如果你的多个步骤之间没有依赖关系只是单纯希望并行执行后汇总结果就不要用thenApply串成链而是用CompletableFuture.allOf组合多个任务。这是另一个话题但很多项目把“顺序工作流”和“并行工作流”混在一起导致设计混乱。本文讨论的是有依赖的顺序链所以allOf不在范围内。10. 总结与后续学习方向到这里顺序工作流异步执行的核心内容基本讲完了。我们要记住几个关键点CompletableFuture可以通过supplyAsync开启异步任务通过thenApply/thenCompose串联有依赖关系的步骤天然保证执行顺序。当链上某个任务抛出异常时后续thenApply等阶段不会继续执行这是默认的异常短路机制不是 bug。需要中断整个流程时不要用exceptionally吞掉异常让它继续向下传播或者重新抛出业务异常。需要降级继续时用exceptionally返回降级值但要设计好降级值在下游的语义。需要统一处理成功和失败分支时用handle。只需要记录日志不改变异常传播时用whenComplete。生产环境必须使用自定义线程池避免使用Executors.newFixedThreadPool这种无界队列写法建议使用ThreadPoolExecutor并配置拒绝策略。链路追踪、超时控制和监控是异步工作流落地必须补齐的配套设施。如果你刚刚接触这个概念建议先花半小时把文章中的最小示例在本机运行一遍故意让某个步骤抛异常观察不同异常处理方法的输出差异。这个实验跑通一次你对CompletableFuture的整个执行链就会有直观感觉。下一步可以继续研究的方向包括多个无依赖的异步任务如何并行执行并用allOf聚合结果orTimeout和completeOnTimeout在 Java 9 中的用法如何把顺序工作流抽象成 DSL 配置让非开发人员也能编排如何在 Spring Boot 项目中结合Async注解和自定义线程池实现类似能力如何在异步链路中引入分布式事务和幂等设计保证最终一致性。技术选型上如果你的任务链非常简单确实可以直接用同步代码没必要引入异步如果你已经遇到了接口响应慢、线程阻塞多或者需要把多步骤任务从请求线程中剥离出来的问题那么CompletableFuture的顺序异步工作流就是非常值得掌握的一把钥匙。这篇文章中的示例代码都基于 JDK 自带的能力没有引入任何第三方框架你可以直接复制到一个空 Java 项目中运行。建议收藏备用下次写业务流程时可以对照检查自己的异常处理策略是否合理。
返回列表