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 會自行決定使用多少執行緒(通常等於 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

1
任務
JAVA 25 SELF, 等級 54, 課堂 3
上鎖
數位會計師的任務分配 💼
數位會計師的任務分配 💼
1
任務
JAVA 25 SELF, 等級 54, 課堂 3
上鎖
破解密碼學代碼 🕵️‍♀️
破解密碼學代碼 🕵️‍♀️
留言
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION