1. Strumienie równoległe (parallelStream): prosto i wygodnie
Jeśli pracowałeś już z Stream API, wiesz, jak wygodnie filtrować, przekształcać i zbierać kolekcje. Ważny bonus: każdy strumień można uczynić równoległym jedną linijką – włącz parallel(), a elementy kolekcji zaczną być przetwarzane współbieżnie.
Jest to szczególnie przydatne dla niezależnych operacji na zbiorze plików: policzyć wiersze, znaleźć podciąg, skopiować, skompresować itd.
Przykład: równoległe zliczanie wierszy we wszystkich plikach katalogu
Wariant 1: sekwencyjnie
import java.nio.file.*;
import java.io.IOException;
import java.util.List;
public class LogLineCounter {
public static void main(String[] args) throws IOException {
Path logDir = Paths.get("logs");
long totalLines = 0;
try (DirectoryStream<Path> stream = Files.newDirectoryStream(logDir, "*.log")) {
for (Path file : stream) {
long lines = Files.lines(file).count();
totalLines += lines;
}
}
System.out.println("Łączna liczba wierszy we wszystkich logach: " + totalLines);
}
}
Komentarz: wszystko odbywa się po kolei, jeden plik po drugim. Jeśli plików jest dużo i są duże, trzeba będzie długo czekać.
Wariant 2: równolegle!
import java.nio.file.*;
import java.io.IOException;
import java.util.stream.Stream;
public class LogLineCounterParallel {
public static void main(String[] args) throws IOException {
Path logDir = Paths.get("logs");
try (Stream<Path> files = Files.list(logDir)) {
long totalLines = files
.filter(path -> path.toString().endsWith(".log"))
.parallel() // i to cała magia!
.mapToLong(file -> {
try (Stream<String> lines = Files.lines(file)) {
return lines.count();
} catch (IOException e) {
e.printStackTrace();
return 0L;
}
})
.sum();
System.out.println("Łączna liczba wierszy we wszystkich logach: " + totalLines);
}
}
}
Komentarz: kluczowa linia – .parallel(). Na procesorze wielordzeniowym program z reguły zadziała wyraźnie szybciej.
Jak to działa?
- parallel() zamienia zwykły strumień w równoległy. Pod spodem używany jest wspólny ForkJoinPool (domyślnie liczba wątków równa liczbie rdzeni).
- Każdy plik jest przetwarzany niezależnie, a wyniki agregowane przez operacje terminalne (np. sum()).
- Jeśli plików jest mało – przyspieszenia może nie być; jeśli są ich setki – zysk bywa zwykle zauważalny.
Ważne!
- Strumienie równoległe nie przyspieszają samych operacji I/O; pozwalają wykonywać kilka operacji jednocześnie. Na szybkich nośnikach (SSD) to pomaga, na wolnych (HDD) można natrafić na „wąskie gardło” dysku.
2. ForkJoinPool: „dziel i rządź” w praktyce
ForkJoin to framework do równoległych obliczeń w stylu „dziel i rządź”: duże zadanie rozbijamy na podzadania, wykonujemy je równolegle i łączymy wyniki. Zarządza tym specjalna pula – ForkJoinPool. To właśnie ona jest używana „za kulisami” strumieni równoległych, ale można nią sterować bezpośrednio, aby zyskać większą elastyczność.
Taki model sprawdza się szczególnie w przypadku struktur rekursywnych (drzewa katalogów), dużych tablic danych oraz zadań, które łatwo zdekomponować na niezależne części.
Przykład: rekursywne wyszukiwanie po drzewie katalogów
Znajdźmy wszystkie ".txt" (włącznie z podfolderami) i policzmy łączną liczbę wierszy.
import java.nio.file.*;
import java.util.concurrent.*;
import java.util.*;
import java.io.IOException;
import java.util.stream.Collectors;
public class FolderLineCounter extends RecursiveTask<Long> {
private final Path dir;
public FolderLineCounter(Path dir) {
this.dir = dir;
}
@Override
protected Long compute() {
List<FolderLineCounter> subTasks = new ArrayList<>();
long lines = 0;
try (DirectoryStream<Path> stream = Files.newDirectoryStream(dir)) {
for (Path entry : stream) {
if (Files.isDirectory(entry)) {
FolderLineCounter task = new FolderLineCounter(entry);
task.fork(); // Uruchamiamy podzadanie
subTasks.add(task);
} else if (entry.toString().endsWith(".txt")) {
try (Stream<String> fileLines = Files.lines(entry)) {
lines += fileLines.count();
} catch (IOException e) {
e.printStackTrace();
}
}
}
} catch (IOException e) {
e.printStackTrace();
}
// Zbieramy wyniki z podzadań
for (FolderLineCounter task : subTasks) {
lines += task.join();
}
return lines;
}
public static void main(String[] args) {
Path root = Paths.get("big_folder");
ForkJoinPool pool = new ForkJoinPool();
FolderLineCounter counter = new FolderLineCounter(root);
long totalLines = pool.invoke(counter);
System.out.println("Łączna liczba wierszy we wszystkich plikach .txt: " + totalLines);
}
}
Co się tutaj dzieje:
- Dla każdego folderu tworzone jest osobne zadanie (FolderLineCounter), dla podkatalogów – własne podzadania (fork()).
- Pliki są zliczane na miejscu, a wyniki sumowane po join() wszystkich podzadań.
Jakie są zalety ForkJoin?
- Skutecznie działa z dużymi hierarchiami (drzewa katalogów).
- Maksymalnie wykorzystuje rdzenie procesora.
- Pozwala precyzyjnie kontrolować równoleglenie i granice zadań.
3. Praktyczne scenariusze zastosowań
Masowa obróbka plików
Na przykład trzeba skopiować tysiące zdjęć do folderu kopii zapasowej.
import java.nio.file.*;
import java.util.List;
import java.util.stream.Collectors;
public class ParallelFileCopier {
public static void main(String[] args) throws Exception {
Path sourceDir = Paths.get("photos");
Path destDir = Paths.get("photos_backup");
Files.createDirectories(destDir);
List<Path> files = Files.list(sourceDir)
.filter(Files::isRegularFile)
.collect(Collectors.toList());
files.parallelStream().forEach(file -> {
try {
Path destFile = destDir.resolve(file.getFileName());
Files.copy(file, destFile, StandardCopyOption.REPLACE_EXISTING);
} catch (Exception e) {
e.printStackTrace();
}
});
System.out.println("Wszystkie pliki zostały skopiowane!");
}
}
Komentarz: każdy plik kopiowany jest w osobnym wątku. Przy dużej liczbie plików przyspieszenie jest zauważalne.
Równoległa kompresja/rozpakowywanie
Analogicznie można zrówoleglić kompresję, przeliczanie hashy, konwersję formatów obrazów itp. przez parallelStream() lub własny ForkJoinPool.
4. Ważne uwagi i ograniczenia
- Operacje I/O nie zawsze zyskują na równoległości. Jeśli dysk lub sieć to „wąskie gardło”, setka zadań równoległych tylko zwiększy konkurencję o zasób.
- Nie uruchamiaj zbyt wielu wątków. Domyślnie strumienie równoległe używają wspólnej puli ForkJoinPool.commonPool() z poziomem równoległości ≈ liczbie rdzeni. Można to zmienić przez właściwość "java.util.concurrent.ForkJoinPool.common.parallelism", ale rób to świadomie.
- Nie zapominaj o synchronizacji. Jeśli kilka wątków zapisuje do tego samego pliku/obiektu – używaj synchronizacji i kolejek; dla niezależnych plików synchronizacja nie jest potrzebna.
5. Krótkie wprowadzenie do FileChannel i dostępu pozycyjnego
Dla bardziej zaawansowanych scenariuszy (np. równoległy odczyt różnych części jednego dużego pliku) użyj java.nio.channels.FileChannel, który wspiera odczyt/zapis pozycyjny.
Przykład: odczyt różnych części pliku w różnych wątkach
import java.nio.channels.FileChannel;
import java.nio.file.*;
import java.nio.ByteBuffer;
public class FileChunkReader implements Runnable {
private final Path path;
private final long position;
private final int size;
public FileChunkReader(Path path, long position, int size) {
this.path = path;
this.position = position;
this.size = size;
}
@Override
public void run() {
try (FileChannel channel = FileChannel.open(path, StandardOpenOption.READ)) {
ByteBuffer buffer = ByteBuffer.allocate(size);
channel.read(buffer, position);
// Przetwarzanie danych
System.out.println("Przeczytano " + buffer.position() + " bajtów od pozycji " + position);
} catch (Exception e) {
e.printStackTrace();
}
}
}
Komentarz: uruchom kilka takich zadań – każde czyta swój zakres. Ale ostrożnie: nie wszystkie dyski i systemy plików lubią silne równoległe obciążenie.
6. Typowe błędy przy równoległym przetwarzaniu plików
Błąd nr 1: równoległy zapis do jednego pliku bez synchronizacji. Dane mieszają się i ulegają uszkodzeniu. Używaj kolejek, buforowania i synchronizacji zapisu.
Błąd nr 2: za dużo równoległości. Strumienie równoległe na słabym sprzęcie/małych plikach generują narzut i mogą spowolnić wykonanie.
Błąd nr 3: ignorowanie błędów I/O. W strumieniach równoległych wyjątki łatwo „zgubić” – obsługuj je wewnątrz lambd, loguj i uwzględniaj awarie.
Błąd nr 4: niezamknięte zasoby. Zawsze używaj try-with-resources dla strumieni/kanałów, w przeciwnym razie doprowadzisz do wycieków i dziwnych błędów.
Błąd nr 5: oczekiwanie „magii” od parallel(). Równoległość przyspiesza tylko przy wystarczającym nakładzie pracy i dostępnych zasobach (CPU, szybki dysk). Samo wywołanie parallel() – to nie jest srebrna kula.
GO TO FULL VERSION