Java中如何实现串行与并行任务混合执行及结果后续串行处理
如何在Java中实现混合串行/并行的任务编排
核心思路
用Java的CompletableFuture实现任务编排是最简洁通用的方案:
- 先把需要串行执行的任务列表封装成单个异步任务单元,确保列表内任务按顺序执行、结果可传递;
- 把所有需要并行的任务单元(包括封装好的串行链、独立任务)提交异步执行,等待全部完成;
- 收集所有前置任务的结果后,再启动下一组串行任务链,传入合并后的结果执行。
通用实现方案
1. 封装串行任务链工具方法
编写两个工具方法,分别处理带结果传递和无返回值的串行任务列表,让任意长度的串行任务都能快速封装成一个CompletableFuture:
import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.function.Function; // 封装带结果传递的串行任务链:前一个任务的输出作为下一个任务的输入 public static <T> CompletableFuture<T> runSerialTasks(List<Function<T, T>> tasks, T initialValue) { CompletableFuture<T> currentFuture = CompletableFuture.completedFuture(initialValue); for (Function<T, T> task : tasks) { currentFuture = currentFuture.thenApply(task); } return currentFuture; } // 封装无返回值的串行任务链:仅按顺序执行,不传递结果 public static CompletableFuture<Void> runSerialRunnables(List<Runnable> tasks) { CompletableFuture<Void> currentFuture = CompletableFuture.completedFuture(null); for (Runnable task : tasks) { currentFuture = currentFuture.thenRun(task); } return currentFuture; }
2. 主流程实现(适配任意数量的并行/串行任务)
以下是针对示例场景的实现,同时支持扩展到任意数量的并行任务组:
import java.util.Arrays; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.function.Function; public class TaskOrchestrator { public static void main(String[] args) throws Exception { // -------------------------- // 1. 定义业务任务(可替换为实际业务逻辑) // -------------------------- // 第一组串行任务:t1→t2→t3 List<Function<String, String>> serialGroup1 = Arrays.asList( input -> { System.out.println("执行任务t1"); return input + "-t1结果"; }, input -> { System.out.println("执行任务t2"); return input + "-t2结果"; }, input -> { System.out.println("执行任务t3"); return input + "-t3结果"; } ); // 并行独立任务:t4(也可以是另一组串行任务链) Function<String, String> parallelTask = input -> { System.out.println("执行任务t4"); try { Thread.sleep(1000); } catch (InterruptedException e) {} // 模拟耗时操作 return input + "-t4结果"; }; // 后续串行任务:t5→t6(接收前置所有任务的合并结果) List<Function<List<String>, String>> postSerialGroup = Arrays.asList( results -> { System.out.println("执行任务t5,合并前置结果:" + results); return String.join(" | ", results); }, mergedResult -> { System.out.println("执行任务t6,处理合并结果:" + mergedResult); return mergedResult + "-t6最终结果"; } ); // -------------------------- // 2. 编排任务执行流程 // -------------------------- // 把第一组串行任务封装为异步Future CompletableFuture<String> serialFuture = runSerialTasks(serialGroup1, "初始参数"); // 把并行任务封装为异步Future(如需自定义线程池,可传入第二个Executor参数) CompletableFuture<String> parallelFuture = CompletableFuture.supplyAsync(() -> parallelTask.apply("初始参数")); // 等待所有前置任务完成,收集结果 CompletableFuture<List<String>> allPreResults = CompletableFuture.allOf(serialFuture, parallelFuture) .thenApply(v -> Arrays.asList(serialFuture.join(), parallelFuture.join())); // 执行后续串行任务链,传入合并后的结果 CompletableFuture<String> finalResult = runSerialTasks(postSerialGroup, allPreResults.join()); // 等待最终任务完成,输出结果 System.out.println("所有任务执行完成,最终结果:" + finalResult.join()); } // 复用工具方法 public static <T> CompletableFuture<T> runSerialTasks(List<Function<T, T>> tasks, T initialValue) { CompletableFuture<T> currentFuture = CompletableFuture.completedFuture(initialValue); for (Function<T, T> task : tasks) { currentFuture = currentFuture.thenApply(task); } return currentFuture; } }
3. 扩展说明
- 支持任意数量的并行任务:如果需要多组串行任务并行执行,只需把每组串行任务都封装为
CompletableFuture,然后将所有Future传入CompletableFuture.allOf()即可; - 自定义线程池:默认使用
ForkJoinPool.commonPool(),如果需要控制线程数量,可创建自定义线程池并传入supplyAsync()、thenApplyAsync()等方法; - 异常处理:可通过
exceptionally()、handle()方法给每个Future添加异常处理逻辑,避免单个任务失败导致整个流程中断。
内容的提问来源于stack exchange,提问作者S2K
相关产品推荐
相关产品推荐

