ARTICLE DETAIL

资讯详情

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

顺序工作流异步执行:CompletableFuture异常传播机制详解

顺序工作流异步执行:CompletableFuture异常传播机制详解 大家好我是你们的老朋友。今天这篇文章我们来聊一个在业务系统里经常遇到、但很多同学没有系统性梳理过的技术组合顺序工作流的异步执行以及在这个过程中非常容易踩坑的CompletableFuture 异常传播问题。很多同学看到“顺序工作流”第一反应是顺序不就是一个个执行吗那不是同步代码 for 循环就行了为什么还要异步其实这里的“顺序”指的是业务依赖上的顺序而“异步执行”指的是线程模型上的非阻塞。两者结合既保证了业务链路的有序性又提升了系统的吞吐和资源利用率。本文会从一个典型的订单处理链路出发结合CompletableFuture的thenApply、thenCompose、exceptionally等核心方法完整演示如何实现一套顺序工作流异步执行框架并重点分析一个高频问题CompletableFuture 出现异常后为什么后面的异步任务不再执行如何控制这种中断行为无论你是刚接触 Java 并发编程的新手还是在项目中已经使用过CompletableFuture但总在异常处理和线程池配置上翻车的同学这篇文章都可以作为一份实战参考。1. 顺序工作流与异步执行先理解问题1.1 什么是顺序工作流工作流Workflow这个概念很宽泛但落到开发层面我们通常说的是一项业务操作被拆分成多个有依赖关系的子任务这些子任务必须按照特定顺序依次完成。举个例子一个简化的订单创建流程可能包含以下步骤创建订单记录。校验并扣减库存。生成物流单号。发送通知消息。这四步之间有明确的先后关系只有订单创建成功了才能去扣库存只有库存扣减成功了才应该生成物流单物流单生成后才可以发送通知。这就是一个典型的顺序工作流。如果用专业术语描述这就是一个**有向无环图DAG**的最简单形态——一条直线每个节点依赖前一个节点的执行结果。1.2 为什么需要异步执行看到这里可能有同学会问既然有顺序依赖那我直接写同步代码不行吗// 同步方式 Order order createOrder(); Stock stock deductStock(order); Logistics logistics createLogistics(stock); sendNotify(logistics);这段同步代码能跑吗能跑。但问题在于线程阻塞每个环节都在阻塞当前线程整体耗时是各环节耗时之和。资源利用率低如果其中一个环节是远程 RPC 调用线程在等待 IO 返回时完全被浪费。扩展性差当工作流节点增多或者需要将部分节点并行执行时同步模型很难优雅扩展。异步执行的核心价值不是“让流程变快”而是释放线程资源提高系统在高并发场景下的吞吐能力。同时异步模型并没有破坏业务上的顺序依赖而是通过回调机制在结果就绪后继续执行下一步。1.3 CompletableFuture 在其中的角色Java 8 引入的CompletableFuture是目前 Java 原生异步编程中最实用的工具类之一。它解决的问题恰好就是工作流编排通过thenApply、thenAccept等实现串行依赖。通过thenCombine、allOf等实现并行聚合。通过exceptionally、handle等实现异常恢复。通过whenComplete实现结果回调。简单来说CompletableFuture把“任务的执行”和“任务之间的依赖关系”解耦了。我们只需要声明依赖关系底层线程池会自动调度执行。2. 环境准备与项目结构2.1 环境说明本文示例代码基于以下环境编写但核心思路不受版本限制组件说明JDK8示例使用 JDK 11 验证构建工具Maven 3.6IDEIntelliJ IDEA 或 Eclipse操作系统Windows / macOS / Linux 均可需要说明的是CompletableFuture在 Java 8 中已经提供但到了 Java 9 之后增加了少量超时相关方法比如orTimeout、completeOnTimeout。如果你的项目跑在 Java 8 上超时控制需要自己通过get(timeout, unit)实现。本文示例以 JDK 11 为主会同时给出兼容写法说明。2.2 创建 Maven 工程创建项目时不需要任何第三方依赖纯 JDK 原生 API 即可。?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 artifactIdworkflow-async-demo/artifactId version1.0.0/version packagingjar/packaging properties maven.compiler.source11/maven.compiler.source maven.compiler.target11/maven.compiler.target /properties dependencies dependency groupIdorg.junit.jupiter/groupId artifactIdjunit-jupiter/artifactId version5.9.2/version scopetest/scope /dependency /dependencies /project这里引入 JUnit 5 只是为了写测试用例不涉及额外框架。2.3 线程池配置使用CompletableFuture时有一个容易被忽略的关键点如果没有显式传入线程池默认使用ForkJoinPool.commonPool()。commonPool是全局共享的默认线程数为CPU 核心数 - 1而且它会被同 JVM 中所有使用默认线程池的异步任务共用。一旦某个任务出现阻塞会直接拖累整个 JVM 的异步任务执行效率。所以在生产项目中强烈建议自定义线程池。示例中我们定义一个专用的线程池配置类package com.example.workflow.config; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ThreadFactory; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; /** * 工作流专用线程池 */ public class WorkflowThreadPool { private static final int CPU_COUNT Runtime.getRuntime().availableProcessors(); private static final ThreadPoolExecutor EXECUTOR new ThreadPoolExecutor( CPU_COUNT * 2, CPU_COUNT * 4, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue(1000), new WorkflowThreadFactory(), new ThreadPoolExecutor.CallerRunsPolicy() ); /** * 自定义线程工厂方便排查问题时识别线程归属 */ static class WorkflowThreadFactory implements ThreadFactory { private final AtomicInteger threadNumber new AtomicInteger(1); Override public Thread newThread(Runnable r) { Thread thread new Thread(r, workflow-thread- threadNumber.getAndIncrement()); thread.setDaemon(false); return thread; } } public static ThreadPoolExecutor getExecutor() { return EXECUTOR; } /** * 优雅关闭线程池 */ public static void shutdown() { EXECUTOR.shutdown(); } }关于参数说明corePoolSize设为CPU_COUNT * 2适合 IO 密集型的业务。maximumPoolSize设为CPU_COUNT * 4留出一定的弹性空间。workQueue使用ArrayBlockingQueue(1000)避免无界队列导致内存溢出。拒绝策略使用CallerRunsPolicy当任务队列满时由提交任务的线程自己执行。这在业务场景中比直接抛出RejectedExecutionException更友好但要注意可能造成提交线程阻塞。3. 核心原理CompletableFuture 的顺序编排与异常传播3.1 常用创建方式CompletableFuture常见的创建方式有三种// 方式一已完成的结果 CompletableFutureString completed CompletableFuture.completedFuture(done); // 方式二异步执行一个 Runnable CompletableFutureVoid runAsyncFuture CompletableFuture.runAsync(() - { System.out.println(执行没有返回值的任务); }, WorkflowThreadPool.getExecutor()); // 方式三异步执行一个 Supplier获取返回值 CompletableFutureString supplyAsyncFuture CompletableFuture.supplyAsync(() - { return 执行有返回值的任务; }, WorkflowThreadPool.getExecutor());注意runAsync适用于没有返回值的任务supplyAsync适用于有返回值的任务。在工作流场景中因为节点之间通常需要传递数据所以supplyAsync更常用。3.2 顺序编排的核心方法CompletableFuture提供了多个方法用于串联任务它们之间的区别非常重要方法返回类型是否接收上个阶段的结果是否返回下个阶段需要的结果适用场景thenApplyCompletableFutureU是是转换结果继续传给下个阶段thenAcceptCompletableFutureVoid是否消费结果不需要继续传递thenRunCompletableFutureVoid否否只关心上个阶段执行完成thenComposeCompletableFutureU是是返回内部 Future上个阶段返回一个 Future需要扁平化这里最核心的是thenApply和thenCompose的区别。thenApply的回调函数返回一个普通对象框架自动包装成新的CompletableFuture。而thenCompose的回调函数本身返回一个CompletableFuture用于避免嵌套的CompletableFutureCompletableFutureT。// thenApply返回普通值 CompletableFutureString future1 CompletableFuture .supplyAsync(() - hello, executor) .thenApply(s - s world); // thenCompose返回 CompletableFuture CompletableFutureString future2 CompletableFuture .supplyAsync(() - hello, executor) .thenCompose(s - CompletableFuture.supplyAsync(() - s world, executor));在工作流编排中如果每个节点的方法本身就是异步的返回CompletableFuture就应该使用thenCompose。如果节点方法是同步的用thenApply就够了。3.3 异常传播机制为什么异常后不再执行后续任务这是本文的重点也是标题中提到的热词核心问题。看下面这段代码package com.example.workflow.demo; import java.util.concurrent.CompletableFuture; public class ExceptionPropagationDemo { public static void main(String[] args) { CompletableFuture.supplyAsync(() - { System.out.println(步骤1创建订单); return ORDER_001; }).thenApply(orderId - { System.out.println(步骤2扣减库存处理订单 orderId); // 模拟这里的异常 throw new RuntimeException(库存不足); }).thenApply(orderId - { System.out.println(步骤3生成物流单处理订单 orderId); return LOGISTICS_001; }).thenAccept(logisticsId - { System.out.println(步骤4发送通知物流单 logisticsId); }).exceptionally(ex - { System.out.println(捕获异常 ex.getMessage()); return null; }); // 等待异步任务执行完成方便观察输出 try { Thread.sleep(2000); } catch (InterruptedException e) { e.printStackTrace(); } } }运行结果步骤1创建订单 步骤2扣减库存处理订单 ORDER_001 捕获异常java.lang.RuntimeException: 库存不足注意步骤 3 和步骤 4 都没有执行。这是CompletableFuture的设计特性异常会沿着依赖链向下传播中断后续所有阶段的执行直到遇到异常处理方法。每个CompletableFuture内部维护了一个状态可能的状态包括已完成正常返回结果。异常完成抛出了异常。未完成。当某个阶段异常完成后它后续依赖它的阶段会直接感知到这个异常并立即以异常状态结束不会再执行该阶段定义的回调函数。所以“异常后不再执行后续异步任务”不是 bug而是框架设计的fail-fast快速失败机制。它避免了在错误数据之上继续执行无意义的操作。比如库存扣减失败了就没有必要继续生成物流单。3.4 异常恢复与兜底那如果某个节点失败后我们不想中断整个流程而是希望走一个兜底逻辑该怎么做有三种方式方法说明exceptionally只有异常时执行返回兜底值handle无论正常还是异常都执行可统一处理whenComplete无论正常还是异常都执行但不改变结果看示例CompletableFuture.supplyAsync(() - { if (true) { throw new RuntimeException(扣减库存失败); } return 100; }, WorkflowThreadPool.getExecutor()) .exceptionally(ex - { System.out.println(库存扣减失败走兜底返回 0); return 0; }) .thenApply(stock - { System.out.println(当前库存 stock); return stock; });使用exceptionally后异常被拦截并转换为一个兜底值后续的thenApply会正常执行。但要注意异常恢复的位置决定了后续节点能否继续执行。如果在exceptionally之后的节点又抛异常异常依然会继续向后传播。4. 完整实战订单处理顺序工作流异步执行下面我们把前面的概念串起来实现一个完整的订单顺序工作流。4.1 业务模型设计我们模拟一个订单处理链路包含四个顺序节点创建订单接收原始请求生成订单 ID。扣减库存根据订单 ID 扣减库存返回扣减数量。生成物流单根据订单信息生成物流单号。发送通知通知用户订单已创建。每个节点都模拟一定的耗时操作比如 RPC 调用、数据库写入使用Thread.sleep代替。整体流程可以概括为请求 → 创建订单 → 扣减库存 → 生成物流单 → 发送通知 → 完成4.2 定义工作流上下文为了保证节点之间能传递数据我们定义一个上下文对象OrderContext按顺序累积每个节点的结果。package com.example.workflow.model; /** * 订单工作流上下文 * 在每个节点之间传递数据 */ public class OrderContext { private String requestId; private String orderId; private Integer stockCount; private String logisticsNo; public String getRequestId() { return requestId; } public void setRequestId(String requestId) { this.requestId requestId; } public String getOrderId() { return orderId; } public void setOrderId(String orderId) { this.orderId orderId; } public Integer getStockCount() { return stockCount; } public void setStockCount(Integer stockCount) { this.stockCount stockCount; } public String getLogisticsNo() { return logisticsNo; } public void setLogisticsNo(String logisticsNo) { this.logisticsNo logisticsNo; } Override public String toString() { return OrderContext{ requestId requestId \ , orderId orderId \ , stockCount stockCount , logisticsNo logisticsNo \ }; } }4.3 定义顺序工作流执行器接下来是关键部分。我们定义一个SequentialWorkflowExecutor类它接收一个初始上下文按顺序注册节点并最终执行整个流程。package com.example.workflow.core; import com.example.workflow.model.OrderContext; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executor; import java.util.function.Function; /** * 顺序工作流执行器 * 使用 CompletableFuture 串联各个节点 */ public class SequentialWorkflowExecutor { /** * 工作流节点定义 */ public interface WorkflowNode extends FunctionOrderContext, OrderContext { } private final Executor executor; private final ListWorkflowNode nodes new ArrayList(); public SequentialWorkflowExecutor(Executor executor) { this.executor executor; } /** * 注册工作流节点 */ public SequentialWorkflowExecutor addNode(WorkflowNode node) { nodes.add(node); return this; } /** * 执行整个工作流 */ public CompletableFutureOrderContext execute(OrderContext initialContext) { CompletableFutureOrderContext future CompletableFuture.completedFuture(initialContext); for (WorkflowNode node : nodes) { future future.thenApplyAsync(node, executor); } return future; } }这里有两个设计点需要解释thenApplyAsync与thenApply的区别thenApply的执行线程取决于上一个任务的执行状态如果上一个任务已经完成则会在当前调用线程中继续执行这可能导致某个阶段的执行线程不确定。thenApplyAsync则强制将任务提交到指定的线程池中执行保证每个节点都在工作流线程池中运行。在实际项目中为了统一线程模型推荐使用thenApplyAsync。节点定义为FunctionOrderContext, OrderContext每个节点接收上一个节点处理后的上下文经过加工后返回新的上下文。由于工作流是顺序执行的这种方式能够天然地传递数据。4.4 实现各个工作流节点接下来创建四个节点类分别对应订单流程的四个步骤。package com.example.workflow.node; import com.example.workflow.core.SequentialWorkflowExecutor.WorkflowNode; import com.example.workflow.model.OrderContext; import java.util.UUID; /** * 节点1创建订单 */ public class CreateOrderNode implements WorkflowNode { Override public OrderContext apply(OrderContext context) { // 模拟耗时操作例如数据库插入 try { Thread.sleep(300); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } String orderId ORDER_ UUID.randomUUID().toString().substring(0, 8).toUpperCase(); context.setOrderId(orderId); System.out.println([创建订单] requestId context.getRequestId() , orderId orderId , thread Thread.currentThread().getName()); return context; } }package com.example.workflow.node; import com.example.workflow.core.SequentialWorkflowExecutor.WorkflowNode; import com.example.workflow.model.OrderContext; /** * 节点2扣减库存 */ public class DeductStockNode implements WorkflowNode { Override public OrderContext apply(OrderContext context) { try { Thread.sleep(400); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } int stockCount 1; context.setStockCount(stockCount); System.out.println([扣减库存] orderId context.getOrderId() , stockCount stockCount , thread Thread.currentThread().getName()); return context; } }package com.example.workflow.node; import com.example.workflow.core.SequentialWorkflowExecutor.WorkflowNode; import com.example.workflow.model.OrderContext; import java.util.UUID; /** * 节点3生成物流单 */ public class CreateLogisticsNode implements WorkflowNode { Override public OrderContext apply(OrderContext context) { try { Thread.sleep(350); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } String logisticsNo LOG_ UUID.randomUUID().toString().substring(0, 8).toUpperCase(); context.setLogisticsNo(logisticsNo); System.out.println([生成物流单] orderId context.getOrderId() , logisticsNo logisticsNo , thread Thread.currentThread().getName()); return context; } }package com.example.workflow.node; import com.example.workflow.core.SequentialWorkflowExecutor.WorkflowNode; import com.example.workflow.model.OrderContext; /** * 节点4发送通知 */ public class SendNotifyNode implements WorkflowNode { Override public OrderContext apply(OrderContext context) { try { Thread.sleep(200); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println([发送通知] orderId context.getOrderId() , logisticsNo context.getLogisticsNo() , thread Thread.currentThread().getName()); return context; } }4.5 运行与验证编写一个启动类来串联整个流程package com.example.workflow; import com.example.workflow.config.WorkflowThreadPool; import com.example.workflow.core.SequentialWorkflowExecutor; import com.example.workflow.model.OrderContext; import com.example.workflow.node.CreateLogisticsNode; import com.example.workflow.node.CreateOrderNode; import com.example.workflow.node.DeductStockNode; import com.example.workflow.node.SendNotifyNode; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; public class OrderWorkflowApplication { public static void main(String[] args) throws Exception { // 1. 构建工作流执行器注册四个顺序节点 SequentialWorkflowExecutor workflowExecutor new SequentialWorkflowExecutor(WorkflowThreadPool.getExecutor()) .addNode(new CreateOrderNode()) .addNode(new DeductStockNode()) .addNode(new CreateLogisticsNode()) .addNode(new SendNotifyNode()); // 2. 构建初始上下文 OrderContext initialContext new OrderContext(); initialContext.setRequestId(REQ_ UUID.randomUUID().toString().substring(0, 8).toUpperCase()); // 3. 执行工作流异步 CompletableFutureOrderContext future workflowExecutor.execute(initialContext); // 4. 设置超时等待 OrderContext result future.get(5, TimeUnit.SECONDS); System.out.println( 工作流执行完成 ); System.out.println(result); // 5. 关闭线程池 WorkflowThreadPool.shutdown(); } }运行结果如下线程名和 ID 每次会不同[创建订单] requestIdREQ_3A2F1B6C, orderIdORDER_7E2F9A1B, threadworkflow-thread-1 [扣减库存] orderIdORDER_7E2F9A1B, stockCount1, threadworkflow-thread-2 [生成物流单] orderIdORDER_7E2F9A1B, logisticsNoLOG_4D8C3E2A, threadworkflow-thread-3 [发送通知] orderIdORDER_7E2F9A1B, logisticsNoLOG_4D8C3E2A, threadworkflow-thread-4 工作流执行完成 OrderContext{requestIdREQ_3A2F1B6C, orderIdORDER_7E2F9A1B, stockCount1, logisticsNoLOG_4D8C3E2A}可以观察到几个细节四个节点按照注册顺序依次执行。每个节点执行在不同的线程上这就是thenApplyAsync将任务提交到线程池的效果。上下文对象在节点之间正确传递每个节点都能拿到前面节点的结果。4.6 异常场景中断后续任务现在我们在DeductStockNode中模拟库存不足的异常package com.example.workflow.node; import com.example.workflow.core.SequentialWorkflowExecutor.WorkflowNode; import com.example.workflow.model.OrderContext; public class DeductStockNode implements WorkflowNode { Override public OrderContext apply(OrderContext context) { try { Thread.sleep(400); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } // 模拟库存不足 throw new RuntimeException(库存不足订单号: context.getOrderId()); } }修改主类在execute后添加异常处理package com.example.workflow; import com.example.workflow.config.WorkflowThreadPool; import com.example.workflow.core.SequentialWorkflowExecutor; import com.example.workflow.model.OrderContext; import com.example.workflow.node.CreateLogisticsNode; import com.example.workflow.node.CreateOrderNode; import com.example.workflow.node.DeductStockNode; import com.example.workflow.node.SendNotifyNode; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; public class OrderWorkflowExceptionApplication { public static void main(String[] args) { SequentialWorkflowExecutor workflowExecutor new SequentialWorkflowExecutor(WorkflowThreadPool.getExecutor()) .addNode(new CreateOrderNode()) .addNode(new DeductStockNode()) .addNode(new CreateLogisticsNode()) .addNode(new SendNotifyNode()); OrderContext initialContext new OrderContext(); initialContext.setRequestId(REQ_ UUID.randomUUID().toString().substring(0, 8).toUpperCase()); CompletableFutureOrderContext future workflowExecutor.execute(initialContext); try { OrderContext result future.get(5, TimeUnit.SECONDS); System.out.println(工作流执行完成: result); } catch (Exception e) { System.out.println(工作流执行失败根因: e.getCause().getMessage()); } finally { WorkflowThreadPool.shutdown(); } } }运行结果[创建订单] requestIdREQ_6D0E8F1A, orderIdORDER_A1B2C3D4, threadworkflow-thread-1 工作流执行失败根因: 库存不足订单号: ORDER_A1B2C3D4可以看到创建订单执行完成后扣减库存抛出异常后续的生成物流单、发送通知节点都没有执行。这正是开头提到的热词问题CompletableFuture 异常后不再执行其他的异步任务。5. 常见问题与排查思路CompletableFuture在工作流编排中踩坑比较多这里整理一份高频问题清单。问题现象常见原因解决思路后续节点不执行但没有异常日志前序节点中异常被吞掉或异常处理位置不当在每个节点出口打印日志检查exceptionally是否返回了null导致空指针被吞异步任务执行在线程上不确定使用了thenApply而不是thenApplyAsync明确指定线程池统一使用thenApplyAsync线程池任务堆积内存飙升有界队列过小或过大拒绝策略不合理根据业务量压测调整队列大小和最大线程数主线程很快就结束了异步任务没执行完没有调用get/join等待或 JVM 提前退出使用get(timeout, unit)阻塞等待或使用CountDownLatch异步任务里的异常没有日志异常被某个上游exceptionally捕获但没打印在全局异常处理点统一打印异常堆栈使用ForkJoinPool.commonPool导致任务互相影响没有传入自定义线程池所有CompletableFuture方法都显式传入线程池get()使用不当导致主线程阻塞过久没有设置超时时间永远使用get(long timeout, TimeUnit unit)代替无参get()依赖链过长代码难维护每个节点直接写thenApply链抽象工作流节点执行器通过注册节点的方式组织流程5.1 异常被吞的典型场景看下面这段代码CompletableFuture.supplyAsync(() - { throw new RuntimeException(业务异常); }, executor) .exceptionally(ex - { // 这里返回 null后续节点拿到 null很可能出现空指针 return null; }) .thenApply(data - { // 这里的 data 是 null业务逻辑再次抛异常 return data.toString(); });这种写法的问题是exceptionally把异常吞掉并返回null后续节点拿到null后发生空指针而空指针又没有合适的处理整个链路相当于“带病运行”排查起来非常困难。建议在exceptionally中如果不能提供真正有业务意义的兜底值就不要返回null而是重新抛出异常或者打印完整堆栈后返回一个标记对象。5.2 排查步骤参考如果线上出现“异步任务没有执行”的问题建议按以下顺序排查确认任务是否真的被提交到了线程池看线程名和任务日志。查看前序节点有没有抛异常异常是否被exceptionally捕获。确认异常分支有没有返回null。检查线程池队列是否已满有没有触发拒绝策略。检查主线程是否还在运行JVM 是否提前退出。检查是否有循环依赖或死锁。6. 最佳实践与工程建议6.1 线程池务必独立配置这是最重要的一条。CompletableFuture默认使用的是ForkJoinPool.commonPool()它是 JVM 全局共享的一旦被长耗时任务阻塞会拖累所有使用默认线程池的异步操作。生产环境中每个业务方向都应该有独立的线程池并且线程池参数要经过压测确认。不要让线上工作流和定时任务、消息消费共用同一个线程池。6.2 每个节点要有超时保护CompletableFuture链上的任何一个节点如果发生阻塞整个链路的get()都会被拖住。因此每个节点内部如果调用了外部接口、数据库、缓存一定要设置超时时间。Java 9 提供了orTimeoutCompletableFuture.supplyAsync(() - { // 模拟慢接口 try { Thread.sleep(3000); } catch (InterruptedException e) { e.printStackTrace(); } return OK; }, executor) .orTimeout(1, TimeUnit.SECONDS) .whenComplete((result, ex) - { if (ex ! null) { System.out.println(任务超时或异常: ex.getMessage()); } });Java 8 环境则通过get(timeout, unit)控制整体超时try { String result future.get(1, TimeUnit.SECONDS); } catch (TimeoutException e) { System.out.println(任务执行超时); // 注意future 本身并不会取消这里只是放弃了等待 }6.3 链路追踪与日志异步执行最大的排查难点在于同一个业务请求可能跨越多个线程日志散落在各个线程中很难串起来。解决办法是使用 TraceId 贯穿始终。可以在工作流上下文中放置一个traceId在每个节点开始和结束时打印System.out.println([创建订单] traceId context.getTraceId() , startTime System.currentTimeMillis() , thread Thread.currentThread().getName());如果公司有完整的链路追踪系统SkyWalking、Zipkin 等优先接入如果没有至少要在业务日志中手动打印 traceId。6.4 异常处理要“分层”建议把异常处理分为两层节点内层每个节点只处理自己能恢复的异常。比如扣库存时发现商品不存在可以在节点内部走兜底逻辑返回一个包含错误标记的上下文。流程外层工作流执行器最外层统一处理“不可恢复”的异常记录错误日志、发送告警、更新流程状态。示例CompletableFutureOrderContext future workflowExecutor.execute(initialContext); future.whenComplete((result, ex) - { if (ex ! null) { // 最外层异常兜底 System.out.println([工作流] 执行异常: ex.getMessage()); // 记录日志、发送告警、更新订单状态为失败 } else { System.out.println([工作流] 执行成功: result); } });6.5 注意异步线程中的上下文传递在实际项目中经常会遇到需要在线程池任务中传递 ThreadLocal 数据的场景比如 Spring Security 的认证信息、SkyWalking 的 TraceId、Dubbo 的隐式传参。默认情况下线程池中的线程无法获取提交任务线程的 ThreadLocal 值。解决方案有三种使用TransmittableThreadLocal阿里开源。将需要传递的数据放入CompletableFuture的入参对象中显式传递。自定义TaskDecoratorSpring 场景下包装 Runnable。对工作流场景建议优先使用第三种思路将上下文对象作为节点入参显式传递尽量避免依赖 ThreadLocal。6.6 优雅关闭线程池应用停机时如果直接 kill 进程正在执行的异步任务会被中断。正确的关闭方式是// 先停止接收新任务 executor.shutdown(); // 等待已提交任务执行完成最多等待 30 秒 boolean terminated executor.awaitTermination(30, TimeUnit.SECONDS); if (!terminated) { // 超时后强制取消剩余任务 executor.shutdownNow(); }6.7 监控指标线程池和工作流都应该有基础监控指标指标说明工作流执行总数统计调用量工作流成功/失败数统计成功率每个节点的平均耗时和 P99 耗时定位性能瓶颈线程池活跃线程数、队列堆积数判断线程池是否饱和拒绝执行次数判断容量是否足够这些指标可以通过 Micrometer、Prometheus 等组件接入也可以简单用日志统计重要的是让问题“可观测”。7. 总结与后续学习方向这篇文章我们围绕“顺序工作流异步执行”这一主题完整走了一遍从概念到落地的过程先理解顺序工作流和异步执行各自的含义再通过CompletableFuture的thenApplyAsync串联四个业务节点最后重点分析了异常传播机制理清了为什么某个节点抛出异常后后续节点不会继续执行。核心收获可以归纳为三点第一CompletableFuture的顺序编排并不复杂关键是在thenApply、thenCompose、thenAccept之间选对方法并且始终显式传入线程池。第二异常中断后续任务不是 bug而是 fail-fast 设计。真正需要做的是在合适的位置拦截异常、记录日志、决定是否走兜底逻辑。第三工作流异步化之后可观测性变得更加重要。线程池参数、超时策略、TraceId 日志、异常兜底这些工程细节往往比主流程代码更值得投入精力。如果你接下来想继续深入建议按这个顺序学习CompletionStage接口的完整方法体系理解所有回调方法的语义。allOf和anyOf实现并行工作流以及并行节点如何聚合结果。自研一个支持 DAG有向无环图的通用异步工作流引擎把复杂业务链路抽象成可配置的节点拓扑。如果你在实际开发中遇到过其他CompletableFuture的诡异问题欢迎在评论区交流。本文的示例代码可以直接复制到本地运行动手验证一遍对异常传播机制的理解会更扎实。
返回列表