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 會自行決定使用多少執行緒(通常等於 CPU 核心數)。
- 任務會自動被拆成子任務,每個子任務都可以在不同核心上執行。
- 效能通常比序列化程式碼更好(尤其在大量資料與多核心系統上)。
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