CodeGym /コース /JAVA 25 SELF /ファイルの並列処理: ForkJoin、Parallel Streams

ファイルの並列処理: ForkJoin、Parallel Streams

JAVA 25 SELF
レベル 59 , レッスン 1
使用可能

1. 並列ストリーム (parallelStream): 手軽で便利

すでに Stream API を使ったことがあるなら、コレクションのフィルタ、変換、収集がどれだけ便利かはご存じでしょう。重要なポイントとして、任意のストリームをわずか 1 行で並列化できます。parallel() を有効にすると、コレクションの要素が並行して処理されます。

これは、ファイル群に対する独立した処理、例えば行数の集計、部分文字列の検索、コピー、圧縮などに特に有効です。

例: ディレクトリ内のすべてのファイルの行数を並列に数える

方法 1: 逐次処理

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("全ログの総行数: " + totalLines);
    }
}

注: すべて順番に、1 ファイルずつ処理します。ファイル数が多かったりサイズが大きい場合は、待ち時間が長くなります.

方法 2: 並列処理!

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() // ここが肝心!
                .mapToLong(file -> {
                    try (Stream<String> lines = Files.lines(file)) {
                        return lines.count();
                    } catch (IOException e) {
                        e.printStackTrace();
                        return 0L;
                    }
                })
                .sum();
            System.out.println("全ログの総行数: " + totalLines);
        }
    }
}

注: 要となるのは .parallel() の 1 行です。マルチコア CPU では、通常は目に見えて高速に動作します。

どう動くのか?

  • parallel() は通常のストリームを並列ストリームに切り替えます。内部では共有の ForkJoinPool が使われ、デフォルトのスレッド数はコア数です。
  • 各ファイルは独立に処理され、結果は終端操作(例: sum())で集約されます。
  • ファイルが少ない場合は高速化しないこともありますが、数百単位になると効果が見えるのが一般的です。

重要!

  • 並列ストリーム自体が I/O を速くするわけではありません。複数の操作を同時に実行できるようにするだけです。高速なストレージ(SSD)では効果がありますが、遅いストレージ(HDD)ではディスクがボトルネックになることがあります。

2. ForkJoinPool: 「分割統治」を実践する

ForkJoin は「分割統治」による並列計算のためのフレームワークです。大きなタスクをサブタスクに分割し、並列に実行して結果を統合します。これを管理するのが専用のプール — ForkJoinPool です。これは並列ストリームの「舞台裏」で使われていますが、より柔軟にするために直接制御することもできます。

このモデルは、再帰的な構造(ディレクトリツリー)、大規模なデータ配列、独立した部分に容易に分割できるタスクに特に適しています。

例: ディレクトリツリーの再帰的探索

すべての ".txt" ファイル(入れ子のフォルダーを含む)を見つけ、合計行数を数えます。

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(); // サブタスクを起動
                    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();
        }
        // サブタスクの結果を集約
        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("すべての .txt ファイルの総行数: " + totalLines);
    }
}

ここで何が起きているか:

  • 各フォルダーごとに個別のタスク(FolderLineCounter)を作成し、サブフォルダーにはそれぞれサブタスク(fork())を起動します。
  • ファイルはその場でカウントし、すべてのサブタスクの join() 後に結果を合計します。

ForkJoin の利点は?

  • 大規模な階層構造(ディレクトリツリー)に対して効率的に動作します。
  • CPU コアを最大限に活用できます。
  • 並列化の粒度やタスク境界をきめ細かく制御できます。

3. 実用的な適用シナリオ

大量ファイルの処理

例えば、数千枚の写真をバックアップフォルダーにコピーする場合。

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("すべてのファイルをコピーしました!");
    }
}

注: 各ファイルは個別のスレッドでコピーされます。ファイル数が多いほど効果がはっきり出ます。

並列圧縮/展開

同様に parallelStream() や独自の ForkJoinPool を使って、圧縮、ハッシュの再計算、画像フォーマット変換なども並列化できます。

4. 重要な注意点と制約

  • I/O は常に並列化で速くなるわけではありません。 ディスクやネットワークがボトルネックなら、数多くの並列タスクはリソース争奪を増やすだけです。
  • スレッドを過剰に増やさない。 デフォルトで並列ストリームは共有の ForkJoinPool.commonPool() を使い、並列度はおおむねコア数です。これはプロパティ "java.util.concurrent.ForkJoinPool.common.parallelism" で変更できますが、十分に理解したうえで行いましょう。
  • 同期を忘れない。 複数のスレッドが同じファイル/オブジェクトへ書き込む場合は同期やキューを使用してください。独立したファイルに対しては同期は不要です。

5. FileChannel と位置指定アクセスの簡単な紹介

高度なシナリオ(例: 大きな 1 つのファイルの異なる部分を並列に読み込む)には、位置指定の読み書きをサポートする java.nio.channels.FileChannel を使います。

例: ファイルの異なる領域を複数スレッドで読み込む

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);
            // データの処理
            System.out.println("読み取り済み " + buffer.position() + " バイト、位置 " + position);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

注: このタスクを複数起動すれば、それぞれが自分の範囲を読み込みます。ただし注意してください。すべてのディスクやファイルシステムが強い並列負荷を好むわけではありません。

6. ファイルの並列処理でよくあるミス

ミス1: 同期なしで 1 つのファイルへ並列書き込みを行う。 データが混ざって破損します。キュー、バッファリング、書き込みの同期を使いましょう。

ミス2: 並列度を上げすぎる。 並列ストリームは、非力なハードウェアや小さなファイルではオーバーヘッドが増え、かえって遅くなることがあります。

ミス3: I/O エラーを無視する。 並列ストリームでは例外が見落とされがちです。ラムダ内で処理し、ログに記録し、失敗を考慮しましょう。

ミス4: リソースを閉じない。 ストリーム/チャネルには常に try-with-resources を使い、リークや不可解なエラーを防ぎましょう。

ミス5: 「魔法」を parallel() に期待する。 並列化が速くなるのは、十分な仕事量と利用可能なリソース(CPU、高速なディスク)がある場合だけです。parallel() を呼ぶだけでは銀の弾丸ではありません。

1
タスク
JAVA 25 SELF, レベル 59, レッスン 1
ロック未解除
ForkJoinを使用してディレクトリ内のファイル総サイズを計算する
ForkJoinを使用してディレクトリ内のファイル総サイズを計算する
1
タスク
JAVA 25 SELF, レベル 59, レッスン 1
ロック未解除
1つのファイルの部分の並列処理
1つのファイルの部分の並列処理
コメント
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION