CodeGym /课程 /JAVA 25 SELF /文件系统的并行遍历:Files.walk + parallel() 与 ForkJoin

文件系统的并行遍历:Files.walk + parallel() 与 ForkJoin

JAVA 25 SELF
第 59 级 , 课程 2
可用

1. 问题:如何高效处理目录中的大量文件

在现代应用中,经常需要在某个文件夹及其子目录中处理大量文件。例如:

  • 统计项目中所有 ".java" 文件的总行数。
  • 查找最近一个月内被修改的全部文件。
  • 按某些条件复制或删除文件。

如果文件很少,一个普通循环就够了。但当文件成千上万,尤其是每个文件上要执行“重”操作(读取、解析、分析)时,耗时会显著增加。

问题:如何加速对大量文件的处理?
答案:使用并行——在多个线程中同时处理文件。

2. 遍历文件系统的工具

Files.walk()

在 Java 8+ 中,出现了遍历目录树的便捷方式——方法 Files.walk()(位于包 java.nio.file)。它返回一个 Stream<Path> 流,包含从给定目录开始的所有文件和文件夹。

示例:

import java.nio.file.*;
import java.util.stream.Stream;

Path start = Paths.get("src");
try (Stream<Path> stream = Files.walk(start)) {
    stream.forEach(System.out::println);
}
  • Files.walk(start) — 返回包含所有文件和文件夹(含子目录)的流。
  • 可以指定最大遍历深度:Files.walk(start, 3)

Files.find()

如果需要立即按条件过滤(例如仅 ".java" 文件),请使用 Files.find()

import java.nio.file.*;
import java.util.stream.Stream;

Path start = Paths.get("src");

try (Stream<Path> stream = Files.find(
        start,
        Integer.MAX_VALUE,
        (path, attr) -> path.toString().endsWith(".java"))) {
    stream.forEach(System.out::println);
}
  • Files.find() 接受一个过滤器(BiPredicate<Path, BasicFileAttributes>),其入参为路径和文件属性。

3. 并行处理:parallel() 与 ForkJoinPool

并行流:.parallel()

任何 Stream 都有方法 parallel()。调用后,元素处理将在多个线程中进行。

Files.walk(start)
    .parallel()
    .forEach(path -> processFile(path));

每个文件都将尽可能并行处理,这在执行“重”操作(读取、解析、计算)时尤其高效。

内部如何工作?ForkJoinPool

并行流使用一个公共线程池——ForkJoinPool.commonPool()。这是一个“智能”的线程池,会在线程之间调度任务。

  • 默认线程数 = 可用处理器数量:Runtime.getRuntime().availableProcessors()
  • “fork/join” 并行模型很适合彼此独立的任务——例如对单个文件的处理。

何时使用 .parallel()

  • 当每个文件的处理彼此独立时。
  • 当操作“很重”(占用 CPU 或长时间等待 IO)时。
  • 当文件很多(数百、上千)时。

不应使用并行流:

  • 如果文件很少(并行化的开销可能大于收益)。
  • 如果需要严格顺序或元素之间存在依赖。

4. 替代方案与并行度调优

何时更适合使用 ExecutorService

并行流适合简单场景。但如果需要:

  • 精确控制线程数(对于 IO-bound 任务,线程数往往需要大于核心数)。
  • 管理队列、取消、重试与错误处理。
  • 构建更复杂的任务流水线。

那么请使用 ExecutorService

import java.nio.file.*;
import java.util.concurrent.*;

ExecutorService executor = Executors.newFixedThreadPool(8);
Files.walk(start)
    .filter(Files::isRegularFile)
    .forEach(path -> executor.submit(() -> processFile(path)));
executor.shutdown();

调整 ForkJoinPool

默认公共池的线程数等于处理器数量。可以通过系统属性进行修改(需在第一次使用并行流之前设置):

System.setProperty("java.util.concurrent.ForkJoinPool.common.parallelism", "16");
  • 设置后,所有并行流将最多使用 16 个线程。

CPU-bound vs IO-bound 任务

  • CPU-bound:重度占用处理器(数学运算、解析、压缩)。线程数 ≈ 核心数。
  • IO-bound:大量磁盘/网络等待。通常让线程数多于核心数更有利。

对于 IO-bound 任务,并行流未必最优——通常增加线程数的自建 ExecutorService 表现更好。

5. 示例:并行查找与处理文件

我们使用并行遍历来统计项目中所有 ".java" 文件的总行数。

import java.nio.file.*;
import java.util.stream.*;
import java.io.IOException;

public class LineCounter {
    public static void main(String[] args) throws IOException {
        Path start = Paths.get("src");

        long totalLines = Files.walk(start)
            .parallel() // 并行处理!
            .filter(p -> p.toString().endsWith(".java"))
            .mapToLong(LineCounter::countLines)
            .sum();

        System.out.println("代码总行数: " + totalLines);
    }

    // 统计文件行数的方法
    private static long countLines(Path path) {
        try (Stream<String> lines = Files.lines(path)) {
            return lines.count();
        } catch (IOException e) {
            System.err.println("读取文件出错: " + path);
            return 0;
        }
    }
}

发生了什么:

  • Files.walk(start) — 遍历所有路径。
  • parallel() — 开启并行处理。
  • filter(...) — 仅保留 ".java" 文件。
  • mapToLong(...) — 统计每个文件的行数。
  • sum() — 求和得到结果。

优点:可以利用多个线程,同时代码仍然简洁。

6. 重要细节与常见错误

  • 并非所有任务都能因并行而加速。对于小规模文件集或很快完成的操作,并行化开销可能反而拖慢程序。
  • 务必关闭资源。处理文件时使用 try-with-resources——这样文件描述符不会“泄漏”。例如,在 Files.lines(path)try(...) 中。
  • 嵌套并行。在其它并行任务中再启动并行流(nested parallelism)很少有效,且可能导致性能下降。
  • 副作用。未经同步不要写入共享结构/文件。应尽量对元素执行“纯”操作。

7. 示意图:并行文件遍历如何工作

flowchart TD
    A["Files.walk(start)"] --> B["Stream<Path>"]
    B --> C{".parallel()?"}
    C -- 否 --> D[普通 forEach]
    C -- 是 --> E["并行 forEach (ForkJoinPool)"]
    E --> F[在多个线程中处理文件]

8. 并行处理文件的常见错误

错误 1: 将并行流用于很小的任务——开销可能大于收益.

错误 2: 期望并行流对 IO-bound 任务的加速与 CPU-bound 一样。对 IO 任务,通常需要线程数更大的 ExecutorService

错误 3: 在 lambda 中未处理的异常——如果不处理 IOException,流可能中断,结果可能不完整。

错误 4: 向共享变量或文件写入时出现竞态——请同步访问或避免副作用。

错误 5: 忘记关闭资源——对所有文件操作使用 try-with-resources

错误 6: 在首次使用之后才尝试修改 ForkJoinPool.commonPool()——通过 System.setProperty(...) 进行的设置必须提前完成。

错误 7: 在其他并行流内部再使用并行流——常常导致性能下降。

评论
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION