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. 基础概念:顺序工作流、异步执行与 CompletableFuture
2.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 version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.example</groupId> <artifactId>async-workflow-demo</artifactId> <version>1.0-SNAPSHOT</version> <packaging>jar</packaging> <properties> <maven.compiler.source>11</maven.compiler.source> <maven.compiler.target>11</maven.compiler.target> <project.build.sourceEncoding>UTF-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(); CompletableFuture<String> 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] 校验商品库存 skuId=SKU-1001, 线程=pool-1-thread-1 [lockStock] 锁定库存 skuId=SKU-1001, quantity=2, 线程=pool-1-thread-1 [createOrder] 创建订单 skuId=SKU-1001, quantity=2, 线程=pool-1-thread-1 [sendMessage] 发送订单通知 orderId=ORD-20250101-001, 线程=pool-1-thread-1注意看,四个步骤全部运行在pool-1-thread-1上。这是因为thenApply默认的异步编排策略:如果前一个任务已经完成,那么下一个thenApply可能会在调用线程上执行;如果前一个任务未完成,则会在前一个任务所在的线程上继续执行。在我们的场景中,每个阶段都是连续提交的,前一个阶段没有完成,所以后续阶段复用了同一个线程。
这一点并不需要过度担心,关键是“顺序”被保证了:后面的步骤一定在接收到前一个步骤的返回值之后才开始。
4.3 更灵活的 thenCompose 写法
thenApply的返回值会直接包装成新的CompletableFuture,如果内部还要返回新的异步任务,就会出现CompletableFuture<CompletableFuture<T>>的嵌套结构。为了避免这种嵌套,可以使用thenCompose。
CompletableFuture<String> 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,看看会发生什么。
CompletableFuture<String> 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只在链上出现异常时触发,它接收异常对象,并返回一个降级结果。这个结果会替换掉异常,后续步骤会继续执行。
CompletableFuture<String> 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。通过判断异常对象是否为空,可以决定是返回正常处理结果,还是降级结果。
CompletableFuture<String> 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方法在阶段完成时触发,无论成功失败。它接收结果和异常,但它的返回值类型是CompletableFuture<T>,其中 T 是上游的结果类型。它的特点是:如果传入参数有异常,whenComplete执行完后,异常仍然会继续向下传播,除非你在回调里抛出别的异常。
CompletableFuture<String> 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"; CompletableFuture<String> 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); CompletableFuture<String> 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 CompletableFuture<String> lockStockWithRetry( OrderWorkflowService service, String skuId, int quantity, ExecutorService executor) { CompletableFuture<String> 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+,直接使用orTimeout:
CompletableFuture<String> 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 项目中运行。建议收藏备用,下次写业务流程时,可以对照检查自己的异常处理策略是否合理。
