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。
GO TO FULL VERSION