You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 05:17:49