【问题标题】:Why is Files.list() parallel stream performing so much slower than using Collection.parallelStream()?为什么 Files.list() 并行流的执行速度比使用 Collection.parallelStream() 慢得多?
【发布时间】:2015-12-17 18:23:50
【问题描述】:

以下代码片段是获取目录列表、对每个文件调用提取方法并将生成的药物对象序列化为 xml 的方法的一部分。

try(Stream<Path> paths = Files.list(infoDir)) {
    paths
        .parallel()
        .map(this::extract)
        .forEachOrdered(drug -> {
            try {
                marshaller.write(drug);
            } catch (JAXBException ex) {
                ex.printStackTrace();
            }
        });
}

这是完全相同的代码,执行完全相同的操作,但使用普通的 .list() 调用来获取目录列表并在结果列表上调用 .parallelStream()

Arrays.asList(infoDir.toFile().list())
    .parallelStream()
    .map(f -> infoDir.resolve(f))
    .map(this::extract)
    .forEachOrdered(drug -> {
        try {
            marshaller.write(drug);
        } catch (JAXBException ex) {
            ex.printStackTrace();
    }
});

我的机器是四核 MacBook Pro,Java v 1.8.0_60(内部版本 1.8.0_60-b27)。

我正在处理 ~ 7000 个文件。 3 次运行的平均值:

第一个版本: .parallel():20 秒。没有.parallel():41 秒

第二版: .parallelStream():12 秒。 .stream():41 秒。

考虑到从流中读取并执行所有繁重工作的extract 方法和执行最终写入的write 调用,并行模式下的这8 秒似乎是一个巨大的差异。

【问题讨论】:

  • 没有parallelStream你的代码表现如何?
  • gee.cs.oswego.edu/dl/html/StreamParallelGuidance.html "目前,基于 JDK IO 的 Stream 源(例如 BufferedReader.lines())主要面向顺序使用,在元素到达时一个接一个地处理。"跨度>
  • 目录流事先不知道它的大小,因此无法平衡工作负载拆分。正如您自己所说,是extract 完成了繁重的工作,而不是读取目录,因此将目录单线程完全读取到内存中并对生成的数组执行并行操作,该数组具有已知的大小,很多更高效。
  • 在每个 Stream 上调用 spliterator().characteristics() 返回的标志是值得的。对我来说(在 Windows 7 64 位上),第一个只返回 DISTINCT,而第二个返回 ORDERED|SIZED|IMMUTABLE|SUBSIZED
  • “根本玩得不好”,你真正的意思是“从无限的、延迟生成的源中有效地提取并行性比从像数组或集合这样的物化有限源中提取并行性更难",这反映在库性能中。

标签: java parallel-processing java-8 nio java-stream


【解决方案1】:

问题在于,Stream API 的当前实现以及 IteratorSpliterator 对于未知大小源的当前实现严重地将这些源拆分为并行任务。你很幸运拥有超过 1024 个文件,否则你根本没有并行化的好处。当前的 Stream API 实现考虑了从 Spliterator 返回的 estimateSize() 值。大小未知的IteratorSpliterator 在拆分前返回Long.MAX_VALUE,其后缀也总是返回Long.MAX_VALUE。其拆分策略如下:

  1. 定义当前批量大小。当前的公式是从 1024 个元素开始并以算术方式递增(2048、3072、4096、5120 等),直到达到 MAX_BATCH 大小(即 33554432 个元素)。
  2. 将输入元素(在您的情况下为路径)消耗到数组中,直到达到批量大小或输入用完为止。
  3. 返回一个 ArraySpliterator 迭代创建的数组作为前缀,将自身作为后缀。

假设您有 7000 个文件。 Stream API 要求估计大小,IteratorSpliterator 返回Long.MAX_VALUE。好的,Stream API 要求 IteratorSpliterator 进行拆分,它从底层 DirectoryStream 收集 1024 个元素到数组并拆分为 ArraySpliterator(估计大小为 1024)和自身(估计大小仍为 Long.MAX_VALUE) .由于Long.MAX_VALUE 远远超过 1024,Stream API 决定继续拆分较大的部分,甚至不尝试拆分较小的部分。所以整体的分裂树是这样的:

                     IteratorSpliterator (est. MAX_VALUE elements)
                           |                    |
ArraySpliterator (est. 1024 elements)   IteratorSpliterator (est. MAX_VALUE elements)
                                           |        |
                           /---------------/        |
                           |                        |
ArraySpliterator (est. 2048 elements)   IteratorSpliterator (est. MAX_VALUE elements)
                                           |        |
                           /---------------/        |
                           |                        |
ArraySpliterator (est. 3072 elements)   IteratorSpliterator (est. MAX_VALUE elements)
                                           |        |
                           /---------------/        |
                           |                        |
ArraySpliterator (est. 856 elements)    IteratorSpliterator (est. MAX_VALUE elements)
                                                    |
                                        (split returns null: refuses to split anymore)

所以之后你有五个并行任务要执行:实际上包含 1024、2048、3072、856 和 0 个元素。请注意,即使最后一个块有 0 个元素,它仍然报告它估计有 Long.MAX_VALUE 元素,因此 Stream API 也会将其发送到 ForkJoinPool。不好的是,Stream API 认为进一步拆分前四个任务是没有用的,因为它们的估计大小要小得多。所以你得到的是输入的非常不均匀的分割,它最多使用四个 CPU 内核(即使你有更多)。如果您的每个元素处理任何元素的时间大致相同,那么整个过程将等待大部分(3072 个元素)完成。所以你可能拥有的最大加速是 7000/3072=2.28x。因此,如果顺序处理需要 41 秒,那么并行流将需要大约 41/2.28 = 18 秒(接近您的实际数字)。

您的解决方案完全没问题。请注意,使用Files.list().parallel() 您还将所有输入Path 元素存储在内存中(在ArraySpliterator 对象中)。因此,如果您手动将它们转储到List 中,您将不会浪费更多内存。像 ArrayList(目前由 Collectors.toList() 创建)这样的数组支持列表实现可以毫无问题地均匀拆分,从而提高速度。

为什么这种情况没有优化?当然这不是不可能的问题(尽管实现可能非常棘手)。对于 JDK 开发人员来说,这似乎不是高优先级的问题。邮件列表中对此主题进行了多次讨论。您可以阅读 Paul Sandoz 的消息 here,他在其中介绍了我的优化工作。

【讨论】:

  • 这是一个很好的解释! +1
  • 感谢您的解释。出于兴趣,这部分是如何工作的:“不好的是,Stream API 认为进一步拆分前四个任务是没有用的,因为它们的估计大小要小得多。”
  • @lwpro2 它看到右侧部分的估计大小远大于左侧部分的估计大小,因此它甚至不再尝试拆分左侧部分。并且正确的部分(实际上包含 0 个元素)拒绝拆分(这是意料之中的)。请注意,所有拆分都是在实际元素处理之前完成的,因此管道在决定拆分位置时并不知道最右边的部分实际上是空的。
  • 谢谢@TagirValeev。似乎在代码中,它是基于Long.MAX_VALUE 错误计算的sizeThreshold。它会停止进一步拆分小于该阈值的任何部分。
【解决方案2】:

作为替代方案,您可以使用专为DirectoryStream 量身定制的自定义拆分器:

public class DirectorySpliterator implements Spliterator<Path> {
    Iterator<Path> iterator;
    long est;

    private DirectorySpliterator(Iterator<Path> iterator, long est) {
        this.iterator = iterator;
        this.est = est;
    }

    @Override
    public boolean tryAdvance(Consumer<? super Path> action) {
        if (iterator == null) {
            return false;
        }
        Path path;
        try {
            synchronized (iterator) {
                if (!iterator.hasNext()) {
                    iterator = null;
                    return false;
                }
                path = iterator.next();
            }
        } catch (DirectoryIteratorException e) {
            throw new UncheckedIOException(e.getCause());
        }
        action.accept(path);
        return true;
    }

    @Override
    public Spliterator<Path> trySplit() {
        if (iterator == null || est == 1)
            return null;
        long e = this.est >>> 1;
        this.est -= e;
        return new DirectorySpliterator(iterator, e);
    }

    @Override
    public long estimateSize() {
        return est;
    }

    @Override
    public int characteristics() {
        return DISTINCT | NONNULL;
    }

    public static Stream<Path> list(Path parent) throws IOException {
        DirectoryStream<Path> ds = Files.newDirectoryStream(parent);
        int splitSize = Runtime.getRuntime().availableProcessors() * 8;
        DirectorySpliterator spltr = new DirectorySpliterator(ds.iterator(), splitSize);
        return StreamSupport.stream(spltr, false).onClose(() -> {
            try {
                ds.close();
            } catch (IOException e) {
                throw new UncheckedIOException(e);
            }
        });
    }
}

只需将Files.list 替换为DirectorySpliterator.list,它就会均匀地并行化,而无需任何中间缓冲。这里我们使用DirectoryStream 生成一个没有任何特定顺序的目录列表的事实,因此每个并行线程只会从中获取一个后续条目(以同步方式,因为我们已经有同步 IO 操作,额外的同步具有 next-to-没有开销)。并行顺序每次都会不同(即使使用forEachOrdered),但Files.list()也不保证顺序。

这里唯一重要的部分是要创建多少并行任务。由于在遍历之前我们不知道文件夹中有多少文件,所以最好使用availableProcessors() 作为基础。我创建了大约8 x availableProcessors() 单个任务,这似乎是一个很好的细粒度/粗粒度折衷方案:如果每个元素的处理不均衡,那么拥有比处理器更多的任务将有助于平衡负载。

【讨论】:

  • 请记住,CPU 的数量不需要是 2 的幂,甚至不需要是偶数。 iirc,AMD 曾经制造过三核 CPU,好吧,您可能会明白为什么它们没有起飞,太多软件无法正常平衡它们的工作......
  • @Holger,是的,我想过。我什至在 3 核 AMD CPU 上工作过。它只是一个禁用了一个核心的 4 核机器(要么它未能通过一些工厂测试,要么出于营销原因,以更便宜的价格出售一些 CPU)。在我的代码中解决这个问题并不难,而且子任务的数量是一个粗略的估计。
  • 你能解释一下为什么需要long e = this.est &gt;&gt;&gt; 1; this.est -= e;吗?
【解决方案3】:

另一种解决方法是在流中使用.collect(Collectors.toList()).parallelStream(),例如

try(Stream<Path> paths = Files.list(infoDir)) {
    paths
        .collect(Collectors.toList())
        .parallelStream()
        .map(this::extract)
        .forEachOrdered(drug -> {
            try {
                marshaller.write(drug);
            } catch (JAXBException ex) {
                ex.printStackTrace();
            }
        });
}

有了这个,你不需要调用.map(f -&gt; infoDir.resolve(f)) 并且性能应该类似于你的第二个解决方案。

【讨论】:

    猜你喜欢
    • 2019-12-09
    • 1970-01-01
    • 2011-04-29
    • 2016-10-23
    • 2021-12-10
    • 1970-01-01
    • 2023-04-06
    • 1970-01-01
    • 2017-11-04
    相关资源
    最近更新 更多