CodeGym /Kursy /JAVA 25 SELF /Równoległe przetwarzanie plików: ForkJoin, strumienie rów...

Równoległe przetwarzanie plików: ForkJoin, strumienie równoległe

JAVA 25 SELF
Poziom 59, Lekcja 1
Dostępny

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.

1
Zadanie
JAVA 25 SELF, poziom 59, lekcja 1
Niedostępne
Obliczanie łącznego rozmiaru plików w katalogu przy użyciu ForkJoin
Obliczanie łącznego rozmiaru plików w katalogu przy użyciu ForkJoin
1
Zadanie
JAVA 25 SELF, poziom 59, lekcja 1
Niedostępne
Równoległe przetwarzanie części jednego pliku
Równoległe przetwarzanie części jednego pliku
Komentarze
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION