CodeGym /Kurslar /JAVA 25 SELF /Faylların paralel emalı: ForkJoin, Parallel Streams

Faylların paralel emalı: ForkJoin, Parallel Streams

JAVA 25 SELF
Səviyyə , Dərs
Mövcuddur

1. Paralel axınlar (parallelStream): sadə və rahat

Əgər artıq Stream API ilə işləmisinizsə, kolleksiyaları süzməyin, çevirməyin və toplamağın nə qədər rahat olduğunu bilirsiniz. Vacib üstünlük: istənilən axını cəmi bir sətrlə paralel etmək olar — parallel() aktivləşdirilir və kolleksiya elementləri eyni vaxtda emal olunmağa başlayır.

Bu, xüsusilə fayl dəstləri üzərində asılı olmayan əməliyyatlar üçün əlverişlidir: sətirləri saymaq, alt sətiri tapmaq, kopyalamaq, sıxmaq və s.

Nümunə: qovluqdakı bütün fayllarda sətirlərin paralel sayılması

Variant 1: ardıcıl

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("Bütün loglarda sətirlərin ümumi sayı: " + totalLines);
    }
}

Şərh: hər şey növbə ilə, fayl-fayl işlənir. Fayllar çoxdursa və böyükdürsə, gözləmək çox çəkəcək.

Variant 2: paralel!

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() // bütün sehr bundan ibarətdir!
                .mapToLong(file -> {
                    try (Stream<String> lines = Files.lines(file)) {
                        return lines.count();
                    } catch (IOException e) {
                        e.printStackTrace();
                        return 0L;
                    }
                })
                .sum();
            System.out.println("Bütün loglarda sətirlərin ümumi sayı: " + totalLines);
        }
    }
}

Şərh: açar sətir — .parallel(). Çoxnüvəli prosessorda proqram, adətən, nəzərəçarpacaq dərəcədə daha sürətli işləyəcək.

Bu necə işləyir?

  • parallel() adi axını paralel edir. Pərdə arxasında ümumi ForkJoinPool istifadə olunur (defolt olaraq axınların sayı nüvələrin sayı qədərdir).
  • Hər fayl müstəqil emal olunur, nəticələr terminal əməliyyatları (məsələn, sum()) vasitəsilə toplanır.
  • Fayl azdırsa — sürətlənmə olmaya da bilər; yüzlərlədirsə — adətən qazanc gözə çarpır.

Vacib!

  • Paralel axınlar I/O əməliyyatlarının özünü sürətləndirmir; onlar bir neçə əməliyyatı eyni vaxtda yerinə yetirməyə imkan verir. Sürətli daşıyıcılarda (SSD) bu kömək edir, yavaşıda (HDD) isə diskin “dar boğazına” dirənə bilərsiniz.

2. ForkJoinPool: “böl və hökm sür” praktikada

ForkJoin — paralel hesablamalar üçün “böl və hökm sür” prinsipi ilə işləyən freymvörkdür: böyük tapşırığı alt tapşırıqlara bölürük, onları paralel icra edirik və nəticələri birləşdiririk. Bunu xüsusi hovuz — ForkJoinPool idarə edir. Paralel axınlarda pərdə arxasında məhz o istifadə olunur, lakin daha çevik idarə üçün onunla birbaşa da işləmək olar.

Bu model rekursiv strukturlar (qovluq ağacları), böyük verilənlər massivləri və asanlıqla müstəqil hissələrə parçalanan tapşırıqlar üçün xüsusilə uyğundur.

Nümunə: qovluq ağacında rekursiv axtarış

Bütün ".txt"-fayllarını (daxili qovluqlar daxil olmaqla) tapacağıq və sətirlərin ümumi sayını hesablayacağıq.

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(); // Alt tapşırığı işə salırıq
                    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();
        }
        // Alt tapşırıqlardan nəticələri toplayırıq
        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("Bütün .txt-fayllarda sətirlərin ümumi sayı: " + totalLines);
    }
}

Burada nə baş verir:

  • Hər qovluq üçün ayrıca tapşırıq (FolderLineCounter) yaradılır, alt qovluqlar üçün — öz alt tapşırıqları (fork()).
  • Fayllar yerində sayılır, bütün alt tapşırıqların join() əməliyyatından sonra nəticələr cəmlənir.

ForkJoin-un üstünlüyü nədədir?

  • Böyük iyerarxiyalarla (qovluq ağacları) effektiv işləyir.
  • Prosessor nüvələrini maksimal dərəcədə istifadə edir.
  • Paralelləşdirmə və tapşırıqların sərhədlərini dəqiq idarə etməyə imkan verir.

3. Praktiki istifadə ssenariləri

Faylların kütləvi emalı

Məsələn, minlərlə şəkili ehtiyat qovluğa kopyalamaq lazımdır.

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("Bütün fayllar kopyalandı!");
    }
}

Şərh: hər fayl ayrı axında kopyalanır. Fayl sayı çox olduqda sürətlənmə nəzərəçarpandır.

Paralel sıxma/çıxarma

Oxşar şəkildə sıxmanı, heşlərin yenidən hesablanmasını, şəkil formatlarının çevrilməsini və s.-ni parallelStream() və ya öz ForkJoinPool-unuzla paralelləşdirmək olar.

4. Vacib qeydlər və məhdudiyyətlər

  • I/O əməliyyatları həmişə paralelləşdirmədən faydalanmır. Disk və ya şəbəkə — “dar boğaz”dırsa, yüzlərlə paralel tapşırıq sadəcə resursa rəqabəti artıracaq.
  • Çox sayda axın işə salmayın. Defolt olaraq paralel axınlar ümumi ForkJoinPool.commonPool()-dan istifadə edir və paralellik səviyyəsi ≈ nüvələrin sayı qədərdir. Bunu "java.util.concurrent.ForkJoinPool.common.parallelism" xüsusiyyəti ilə dəyişmək olar, amma bunu düşünülmüş şəkildə edin.
  • Sinxronizasiyanı unutmayın. Bir neçə axın eyni fayla/obyektə yazırsa — sinxronizasiya və növbələrdən istifadə edin; müstəqil fayllar üçün — sinxronizasiya lazım deyil.

5. FileChannel və pozisional girişlə qısa tanışlıq

İrəli səviyyə ssenarilər üçün (məsələn, eyni böyük faylın müxtəlif hissələrini paralel oxumaq) java.nio.channels.FileChannel istifadə edin; o, pozisional oxuma/yazmanı dəstəkləyir.

Nümunə: faylın müxtəlif hissələrinin müxtəlif axınlarda oxunması

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);
            // Məlumatların emalı
            System.out.println("Oxunub " + buffer.position() + " bayt pozisiyadan " + position);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

Şərh: bu cür tapşırıqlardan bir neçəsini başladın — hər biri öz diapazonunu oxuyur. Amma ehtiyatlı olun: bütün disklər və fayl sistemləri güclü paralel yüklənməni sevmir.

6. Faylların paralel emalında tipik səhvlər

Səhv №1: sinxronizasiya olmadan bir fayla paralel yazma. Məlumatlar qarışır və korlanır. Növbələrdən, buferləşdirmədən və yazının sinxronizasiyasından istifadə edin.

Səhv №2: həddindən artıq paralellik. Zəif avadanlıqda/kiçik fayllarda paralel axınlar əlavə xərc yaradır və icranı ləngidə bilər.

Səhv №3: I/O xətalarını görməməzlikdən gəlmək. Paralel axınlarda istisnaları asanlıqla “itirmək” olar — onları lambda daxilində emal edin, loqlaşdırın, nasazlıqları nəzərə alın.

Səhv №4: bağlanmamış resurslar. Həmişə axınlar/kanallar üçün try-with-resources istifadə edin, yoxsa sızmalar və qəribə xətalarla üzləşəcəksiniz.

Səhv №5: parallel()-dən “sehr” gözləmək. Paralellik yalnız işin kifayət qədər həcmi və mövcud resurslar (CPU, sürətli disk) olduqda sürətləndirir. Təkcə parallel() çağırışı — gümüş güllə deyil.

Şərhlər
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION