【发布时间】: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