CodeGym /课程 /JAVA 25 SELF /ForkJoinPool 与 RecursiveTask:递归任务

ForkJoinPool 与 RecursiveTask:递归任务

JAVA 25 SELF
第 54 级 , 课程 3
可用

1. ForkJoinPool:它是什么,有什么用

ForkJoinPool 是一种实现“分而治之”(divide and conquer)方法的专用线程池。它的目标是在可以把大任务拆分为一组相互独立的子任务、并行执行并最终合并结果的场景下尽可能高效地并行化工作。

  • Fork(拆分)——将任务拆分为子任务。
  • Join(合并)——收集子任务结果并汇总。

ForkJoinPool 是 Java 并行流的“心脏”:当你编写 list.parallelStream() 时,内部用的就是它。当然你也可以直接使用它,从而获得更多控制权。

ForkJoinPool 何时特别有用

ForkJoinPool 在易于拆分为独立部分的任务上大放异彩:例如处理超大数组,每个片段独立处理,然后合并结果。

  • 任务易于拆分为相互独立的子任务:排序、搜索、求和等。
  • 子任务在规模上大致相当且彼此独立。
  • 需要利用所有 CPU 内核以获得最大速度。
+---------------------+
|       大型任务      |
+---------------------+
          |
          v
+---------+---------+
|   子任务 1        |
|   子任务 2        |
|   ...            |
+-------------------+
          |
          v
+---------+---------+
|       结果        |
+-------------------+

这正是“分而治之”的工作方式:先拆分——并行计算——再合并。

2. RecursiveTask 和 RecursiveAction:一枚硬币的两面

ForkJoinPool 中,任务通过能拆分子任务并合并结果的特殊类来表达。RecursiveTask<T> 会返回结果,而 RecursiveAction 不返回结果。实践中更常用 RecursiveTask,例如返回和、最大值或计数。

要创建这样的任务,我们继承并实现 compute() 方法。在其中描述逻辑:如果任务很小——直接求解;如果很大——拆分为子任务,通过 fork() 并行启动,并使用 join() 合并结果。这样就形成了自然的递归并行。

3. 语法与示例:并行计算数组之和

假设我们有一个很大的数字数组,希望快速计算所有元素的总和。

步骤 1:任务类

import java.util.concurrent.RecursiveTask;

public class ArraySumTask extends RecursiveTask<Long> {
    private static final int THRESHOLD = 1_000; // 任务拆分阈值
    private final int[] array;
    private final int start, end;

    public ArraySumTask(int[] array, int start, int end) {
        this.array = array;
        this.start = start;
        this.end = end;
    }

    @Override
    protected Long compute() {
        // 如果任务很小——直接计算
        if (end - start <= THRESHOLD) {
            long sum = 0;
            for (int i = start; i < end; i++) {
                sum += array[i];
            }
            return sum;
        } else {
            // 将任务拆分为两个子任务
            int mid = (start + end) / 2;
            ArraySumTask leftTask = new ArraySumTask(array, start, mid);
            ArraySumTask rightTask = new ArraySumTask(array, mid, end);

            // 并行启动子任务
            leftTask.fork(); // 异步
            long rightResult = rightTask.compute(); // 同步
            long leftResult = leftTask.join(); // 等待左侧完成

            // 合并结果
            return leftResult + rightResult;
        }
    }
}
  • 如果任务很小(小于阈值 THRESHOLD)——用普通循环直接求和。
  • 如果任务较大——拆分为两部分,一部分通过 fork() 异步启动,另一部分用 compute() 同步计算,然后通过 join() 合并。

步骤 2:通过 ForkJoinPool 启动任务

import java.util.concurrent.ForkJoinPool;

public class ForkJoinSumDemo {
    public static void main(String[] args) {
        int[] numbers = new int[10_000_000];
        for (int i = 0; i < numbers.length; i++) {
            numbers[i] = 1; // 为简单起见——和应等于数组长度
        }

        ForkJoinPool pool = new ForkJoinPool(); // 默认——等于 CPU 核心数

        ArraySumTask task = new ArraySumTask(numbers, 0, numbers.length);

        long result = pool.invoke(task); // 启动任务

        System.out.println("数组元素之和:" + result);
    }
}

它如何工作?

  • ForkJoinPool 会自行决定使用多少线程(通常等于内核数)。
  • 任务会自动拆分为子任务,每个子任务都可能在不同内核上执行。
  • 性能通常优于顺序代码(尤其在大数据量和多核系统上)。

4. ForkJoinPool 如何工作:一点“底层原理”

Work-Stealing(工作窃取)

ForkJoinPool 实现了“工作窃取”:如果某个线程没有任务了,它会从其他线程处“偷取”工作。这保证了负载的高效均衡与所有内核的充分利用。

基本算法

  • 主任务被拆分为子任务。
  • 子任务被放入专用队列。
  • 线程从自己的队列取任务,当队列为空时去邻居那里“找活”。
  • 当全部完成后,合并结果。

工作示意图

flowchart TD
    A[主任务] --> B1[子任务 1]
    A --> B2[子任务 2]
    B1 --> C1[小任务 1]
    B1 --> C2[小任务 2]
    B2 --> C3[小任务 3]
    B2 --> C4[小任务 4]
    C1 --> D[结果合并]
    C2 --> D
    C3 --> D
    C4 --> D

5. RecursiveAction——当不需要返回结果时

如果只是需要并行执行某些操作而不返回结果,请使用 RecursiveAction。典型示例——并行填充数组、打印、就地排序等。

import java.util.concurrent.RecursiveAction;

public class PrintTask extends RecursiveAction {
    private static final int THRESHOLD = 100;
    private final int[] array;
    private final int start, end;

    public PrintTask(int[] array, int start, int end) {
        this.array = array;
        this.start = start;
        this.end = end;
    }

    @Override
    protected void compute() {
        if (end - start <= THRESHOLD) {
            for (int i = start; i < end; i++) {
                System.out.print(array[i] + " ");
            }
        } else {
            int mid = (start + end) / 2;
            invokeAll(
                new PrintTask(array, start, mid),
                new PrintTask(array, mid, end)
            );
        }
    }
}

6. 实践:并行查找数组中的最大值

import java.util.concurrent.RecursiveTask;

public class MaxFindTask extends RecursiveTask<Integer> {
    private static final int THRESHOLD = 1000;
    private final int[] array;
    private final int start, end;

    public MaxFindTask(int[] array, int start, int end) {
        this.array = array;
        this.start = start;
        this.end = end;
    }

    @Override
    protected Integer compute() {
        if (end - start <= THRESHOLD) {
            int max = array[start];
            for (int i = start + 1; i < end; i++) {
                if (array[i] > max) max = array[i];
            }
            return max;
        } else {
            int mid = (start + end) / 2;
            MaxFindTask left = new MaxFindTask(array, start, mid);
            MaxFindTask right = new MaxFindTask(array, mid, end);
            left.fork();
            int rightResult = right.compute();
            int leftResult = left.join();
            return Math.max(leftResult, rightResult);
        }
    }
}

运行:

import java.util.concurrent.ForkJoinPool;

public class ForkJoinMaxDemo {
    public static void main(String[] args) {
        int[] array = new int[5_000_000];
        for (int i = 0; i < array.length; i++) {
            array[i] = (int)(Math.random() * 1_000_000);
        }

        ForkJoinPool pool = new ForkJoinPool();
        MaxFindTask task = new MaxFindTask(array, 0, array.length);

        int max = pool.invoke(task);

        System.out.println("最大值:" + max);
    }
}

7. ForkJoinPool 的优势与限制

优势

  • 自动负载均衡。 Work-stealing 能高效利用所有内核。
  • 方便。 无需手动创建和管理线程。
  • 高性能。 尤其适用于大任务与多核系统。
  • 灵活。 可以按需把任务拆分为任意数量的部分。

限制

  • 子任务之间耦合度高。 如果子任务经常相互等待,收益会下降。
  • 任务过于细碎。 拆分/同步的开销可能会“吞噬”并行带来的优势。
  • 副作用。 未同步地修改共享变量会导致竞争条件(race condition)。
  • 适用性。 适合能拆分为独立部分的任务。

8. 使用 ForkJoinPool 和 RecursiveTask 时的常见错误

错误 №1:将任务切得过细。 如果阈值(THRESHOLD)过小,会产生大量微小任务——创建与同步的开销将超过并行的收益。请通过实验调整阈值:最佳值往往在数千或数万元素量级。

错误 №2:使用共享的可变变量。 如果子任务未经同步就写入共享变量——会产生数据竞争(race condition)。通过 compute() 返回结果,并只在 join() 中合并。

错误 №3:错误使用 fork/join。 忘记调用 fork()join()——子任务不会并行运行或结果会“丢失”。请仔细关注调用顺序。

错误 №4:在 ForkJoinPool 之外运行 ForkJoinTask。 如果直接调用任务的 compute(),它会在当前线程执行,不会并行。想要真正的并行效果,请使用 pool.invoke()pool.submit()

错误 №5:忽略异常。 如果任务中发生异常,它会在调用 join()invoke() 时显现。不要忘记处理错误。

错误 №6:将 ForkJoinPool 用于包含阻塞的任务。 ForkJoinPool 不适合经常阻塞的任务(等待 I/O 等)。这种情况下更适合使用 ExecutorService

评论
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION