【问题标题】:Make Stream Parallel on the Result of flatMap在 flatMap 的结果上使 Stream Parallel
【发布时间】:2021-01-09 05:45:10
【问题描述】:

考虑以下简单代码:

Stream.of(1)
  .flatMap(x -> IntStream.range(0, 1024).boxed())
  .parallel() // Moving this before flatMap has the same effect because it's just a property of the entire stream
  .forEach(x -> {
     System.out.println("Thread: " + Thread.currentThread().getName());
  });

很长一段时间以来,我都认为即使在flatMap 之后,Java 也会对元素进行并行执行。但是上面的代码打印了所有的“Thread:main”,这证明我的想法是错误的。

在flatMap 之后使其并行的一种简单方法是收集然后再次流式传输:

Stream.of(1)
  .flatMap(x -> IntStream.range(0, 1024).boxed())
  .parallel() // Moving this before flatMap has the same effect because it's just a property of the entire stream
  .collect(Collectors.toList())
  .parallelStream()
  .forEach(x -> {
     System.out.println("Thread: " + Thread.currentThread().getName());
  });

我想知道是否有更好的方法,以及flatMap 的设计选择,它只在调用之前并行化流,而不是在调用之后。

========= 关于问题的更多说明 ========

从一些答案来看,我的问题似乎没有完全传达。正如@Andreas 所说,如果我从 3 个元素的 Stream 开始,可能会有 3 个线程在运行。

但我的问题确实是:根据this post,Java Stream 使用一个常见的 ForkJoinPool,其默认大小等于核心数的 1。现在假设我有 64 个内核,那么我希望我上面的代码在flatMap 之后会看到许多不同的线程,但实际上,它只看到一个(或在 Andreas 的情况下为 3 个)。顺便说一句,我确实使用isParallel 观察到流是并行的。

说实话,我问这个问题并不是出于纯粹的学术兴趣。我在一个项目中遇到了这个问题,该项目提供了用于转换数据集的一长串流操作。链从单个文件开始,并通过flatMap 分解为很多元素。但显然,在我的实验中,它并没有完全利用我的机器(它有 64 个核心),而是只使用了一个核心(通过观察 cpu 使用情况)。

【问题讨论】:

  • 我认为对parallel() 的调用是在flatMap 之前还是之后进行的并不重要。您正在询问规范中没有的内容,即它是特定于实现的。它可能会从一个版本更改为另一个版本,因此不能相信代码在版本之间的行为方式相同。我认为,就目前而言,一切都取决于流的大小。尝试使用更大的值和不同版本的 Java
  • flatMap 业务的意义何在?它只是掩盖了这个问题。
  • 我还要注意,您可以随时致电isParallel 来查看您认为的流实际上是并行的还是顺序的。
  • 这是 OpenJDK 实现的一个已知限制。如果他们要改变这一点,我看到其他一些问题,他们必须首先解决。例如,单个元素流是隐式无序的,现在没有效果,但是当为“子流”启用并行化时,将整个流视为无序可能会导致意外。如果你想执行并行递归文件处理this answer 可以作为一个灵感。
  • @JohnMeyer 内部已在this old Q&A 中讨论。虽然同时缺少短路的问题已得到解决,但仅按顺序迭代子流的基本行为并未改变。

标签: java parallel-processing java-stream


【解决方案1】:

我想知道 [...] flatMap 的设计选择,它只在调用之前并行化流,而不是在调用之后。

你错了。 flatMap 之前和之后的所有步骤都是并行运行的,但它只会在线程之间拆分 original 流。然后flatMap 操作由一个这样的线程处理,并且它的流不会被拆分。

由于您的原始流只有 1 个元素,因此无法拆分,因此 parallel 无效。

尝试更改为Stream.of(1, 2, 3),您将看到forEach,它在flatMap之后,实际上在3个不同的线程中运行。

【讨论】:

  • 这实际上不是指定行为。
  • @chrylis-cautiouslyoptimistic- 我从来没有说过,我正在纠正 OP 的错误声明,即 flatMap 之后的步骤没有并行执行。我已经描述了 observed 行为,与 question 一样。我没有对“特定”行为提出任何声明,并且将来不会改变。
【解决方案2】:

documentation for forEach 指定:

对于任何给定的元素,可以在库选择的任何时间和线程中执行操作。

特别是,“在调用线程上执行所有操作”似乎是一个很好的广泛安全的实现。

请注意,您尝试并行化流并不需要任何特定的并行性,但您更有可能看到这样的效果:

IntStream.range(0, 1024).boxed()
  .parallel()
  .map(i -> "Thread: " + Thread.currentThread().getName())
  .forEach(System.out::println);

【讨论】:

    【解决方案3】:

    对于像我这样迫切需要并行化 flatMap 并需要一些实用解决方案的人来说,不仅仅是历史和理论。对于那些在并行化之前不考虑收集所有项目之间的人。

    我想出的最简单的解决方案是手动进行展平,基本上是将其替换为map + reduce(Stream::concat)。

    我已经在另一个帖子中回答了同样的问题,详情请参阅https://stackoverflow.com/a/66386078/3606820

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2023-04-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-08-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多