1. ForkJoinPool: co to jest i do czego służy
ForkJoinPool — to specjalna pula wątków, realizująca podejście „dziel i zwyciężaj” (divide and conquer). Jej zadaniem jest możliwie najefektywniejsze zrównoleglenie pracy, gdy duże zadanie można rozbić na szereg niezależnych podzadań, wykonać je równolegle, a następnie połączyć wyniki.
- Fork (podzielić) — zadanie dzieli się na podzadania.
- Join (połączyć) — wyniki podzadań są zbierane w rezultat końcowy.
ForkJoinPool — to serce równoległych strumieni w Javie: kiedy piszesz list.parallelStream(), pod spodem używany jest właśnie on. Możesz jednak stosować go także bezpośrednio, zyskując większą kontrolę.
Kiedy ForkJoinPool jest szczególnie przydatny
ForkJoinPool błyszczy w zadaniach, które łatwo podzielić na niezależne części: na przykład przetwarzanie bardzo dużych tablic, gdy każdy fragment jest przetwarzany osobno, a następnie wyniki są łączone.
- Zadanie łatwo rozbić na niezależne podzadania: sortowanie, wyszukiwanie, sumowanie.
- Podzadania są mniej więcej jednakowe pod względem rozmiaru i nie zależą od siebie.
- Należy wykorzystać wszystkie rdzenie procesora dla maksymalnej szybkości.
+---------------------+
| Duże zadanie |
+---------------------+
|
v
+---------+---------+
| Podzadanie 1 |
| Podzadanie 2 |
| ... |
+-------------------+
|
v
+---------+---------+
| Wyniki |
+-------------------+
Właśnie tak działa „dziel i zwyciężaj”: podzieliliśmy — policzyliśmy równolegle — połączyliśmy.
2. RecursiveTask i RecursiveAction: dwie strony tego samego medalu
W ForkJoinPool zadania są realizowane przez specjalne klasy, które potrafią dzielić się na podzadania i łączyć wyniki. RecursiveTask<T> zwraca rezultat, a RecursiveAction — nie. W praktyce częściej używa się RecursiveTask, aby na przykład zwrócić sumę, maksimum lub liczność.
Aby stworzyć takie zadanie, dziedziczymy po klasie i implementujemy metodę compute(). W niej opisujemy logikę: jeśli zadanie jest małe — rozwiązujemy je od razu; jeśli duże — dzielimy na podzadania, uruchamiamy je równolegle przez fork() i łączymy wyniki za pomocą join() wewnątrz compute(). Tak powstaje naturalny rekurencyjny równoległy algorytm.
3. Składnia i przykład: równoległe obliczanie sumy tablicy
Załóżmy, że mamy dużą tablicę liczb i chcemy szybko policzyć sumę wszystkich elementów.
Krok 1. Klasa zadania
import java.util.concurrent.RecursiveTask;
public class ArraySumTask extends RecursiveTask<Long> {
private static final int THRESHOLD = 1_000; // Próg podziału zadania
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() {
// Jeśli zadanie jest małe — liczymy bezpośrednio
if (end - start <= THRESHOLD) {
long sum = 0;
for (int i = start; i < end; i++) {
sum += array[i];
}
return sum;
} else {
// Dzielimy zadanie na dwa podzadania
int mid = (start + end) / 2;
ArraySumTask leftTask = new ArraySumTask(array, start, mid);
ArraySumTask rightTask = new ArraySumTask(array, mid, end);
// Uruchamiamy podzadania równolegle
leftTask.fork(); // Asynchronicznie
long rightResult = rightTask.compute(); // Synchronicznie
long leftResult = leftTask.join(); // Czekamy na zakończenie lewego podzadania
// Łączymy wynik
return leftResult + rightResult;
}
}
}
- Jeśli zadanie jest małe (mniejsze niż próg THRESHOLD) — sumę liczymy zwykłą pętlą.
- Jeśli duże — dzielimy na dwa, jedno uruchamiamy asynchronicznie przez fork(), drugie liczymy synchronicznie przez compute(), następnie łączymy przez join().
Krok 2. Uruchomienie zadania przez 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; // Dla prostoty — suma powinna być równa długości tablicy
}
ForkJoinPool pool = new ForkJoinPool(); // Domyślnie — liczba wątków równa liczbie rdzeni
ArraySumTask task = new ArraySumTask(numbers, 0, numbers.length);
long result = pool.invoke(task); // Uruchomienie zadania
System.out.println("Suma elementów tablicy: " + result);
}
}
Jak to działa?
- ForkJoinPool sam decyduje, ile wątków użyć (zwykle — tyle, ile jest rdzeni).
- Zadanie automatycznie dzieli się na podzadania, z których każde może być wykonywane na osobnym rdzeniu.
- Wydajność zwykle jest wyższa niż w kodzie sekwencyjnym (zwłaszcza na dużych danych i systemach wielordzeniowych).
4. Jak działa ForkJoinPool: trochę „pod maską”
Work-stealing (kradzież pracy)
ForkJoinPool implementuje „kradzież pracy”: jeśli jakiemuś wątkowi skończą się zadania, „kradnie” pracę od innego wątku. Zapewnia to efektywne równoważenie obciążenia i wykorzystanie wszystkich rdzeni.
Podstawowy algorytm
- Główne zadanie dzieli się na podzadania.
- Podzadania są umieszczane w wyspecjalizowanych kolejkach.
- Wątki pobierają zadania ze swoich kolejek, a gdy te się opróżnią — „szukają” pracy u sąsiadów.
- Gdy wszystko zostanie wykonane, wyniki są łączone.
Schemat działania
flowchart TD
A[Zadanie główne] --> B1[Podzadanie 1]
A --> B2[Podzadanie 2]
B1 --> C1[Małe zadanie 1]
B1 --> C2[Małe zadanie 2]
B2 --> C3[Małe zadanie 3]
B2 --> C4[Małe zadanie 4]
C1 --> D[Łączenie wyników]
C2 --> D
C3 --> D
C4 --> D
5. RecursiveAction — jeśli nie trzeba zwracać wyniku
Jeśli trzeba po prostu coś wykonać równolegle i nie zwracać wyniku, użyj RecursiveAction. Typowe przykłady — zrównoleglone wypełnianie tablicy, drukowanie, sortowanie „na miejscu” itp.
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. Praktyka: równoległe wyszukiwanie wartości maksymalnej w tablicy
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);
}
}
}
Uruchomienie:
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("Maksymalna wartość: " + max);
}
}
7. Zalety i ograniczenia ForkJoinPool
Zalety
- Automatyczne równoważenie obciążenia. Work-stealing pozwala efektywnie wykorzystać wszystkie rdzenie.
- Wygoda. Nie trzeba ręcznie tworzyć i zarządzać wątkami.
- Wysoka wydajność. Zwłaszcza w dużych zadaniach i na systemach wielordzeniowych.
- Elastyczność. Można dzielić zadania na tyle części, ile potrzeba.
Ograniczenia
- Silne powiązanie między podzadaniami. Jeśli podzadania często czekają na siebie nawzajem, korzyść maleje.
- Zbyt drobne zadania. Narzuty na podział/synchronizację mogą „zjeść” przewagę.
- Efekty uboczne. Nie wolno zmieniać współdzielonych zmiennych bez synchronizacji — skończy się to race condition.
- Zastosowanie. Dobre dla zadań dzielonych na niezależne części.
8. Typowe błędy przy pracy z ForkJoinPool i RecursiveTask
Błąd nr 1: Zbyt drobny podział zadania. Jeśli próg (THRESHOLD) jest zbyt mały, powstanie wiele malutkich zadań — koszty ich tworzenia i synchronizacji przewyższą zysk z równoległości. Eksperymentuj z progiem: optymalne wartości to często tysiące lub dziesiątki tysięcy elementów.
Błąd nr 2: Używanie wspólnych zmiennych modyfikowalnych. Jeśli podzadania zapisują do wspólnej zmiennej bez synchronizacji — dostaniesz wyścigi danych (race condition). Zwracaj wynik przez compute() i łącz tylko w join().
Błąd nr 3: Nieprawidłowe użycie fork/join. Zapomniano wywołać fork() lub join() — i podzadanie nie uruchomi się równolegle albo wynik „zginie”. Uważnie pilnuj kolejności wywołań.
Błąd nr 4: Uruchamianie ForkJoinTask poza ForkJoinPool. Jeśli po prostu wywołasz compute() na zadaniu, wykona się ono w bieżącym wątku, bez równoległości. Do „prawdziwej magii” używaj pool.invoke() lub pool.submit().
Błąd nr 5: Ignorowanie wyjątków. Jeśli w zadaniu wystąpi wyjątek, ujawni się on przy wywołaniu join() lub invoke(). Nie zapominaj obsługiwać błędów.
Błąd nr 6: Używanie ForkJoinPool do zadań z blokadami. ForkJoinPool słabo nadaje się do zadań, które często się blokują (czekają na I/O itp.). W takich przypadkach lepiej użyć ExecutorService.
GO TO FULL VERSION