had*_*ddr 5 java parallel-processing fork-join java-8 java-stream
我正在使用最新的 Java 8 lambda 和并行流来处理数据。我的代码如下:
ForkJoinPool forkJoinPool = new ForkJoinPool(10);
List<String> files = Arrays.asList(new String[]{"1.txt"});
List<String> result = forkJoinPool.submit(() ->
files.stream().parallel()
.flatMap(x -> stage1(x)) //at this stage we add more elements to the stream
.map(x -> stage2(x))
.map(x -> stage3(x))
.collect(Collectors.toList())
).get();
Run Code Online (Sandbox Code Playgroud)
该流以一个元素开始,但在第二阶段添加更多元素。我的假设是该流应该并行运行,但在这种情况下仅使用一个工作线程。
如果我从 2 个元素开始(即,我将第二个元素添加到初始列表中),则会生成 2 个线程来处理流,依此类推...如果我没有显式地将流提交到 ForkJoinPool,也会发生这种情况。
问题是:它的行为是否记录在案,或者在实施过程中可能会发生变化?有什么方法可以控制这种行为并允许更多线程,无论初始列表如何?
您观察到的是特定于实现的行为,而不是指定的行为。
当前的 JDK 8 实现查看最外层的流Spliterator,并将其用作分割并行工作负载的基础。由于该示例在原始源流中只有一个元素,因此无法拆分,并且该流以单线程运行。这对于返回零个、一个或少数元素的常见(但绝不是唯一)情况非常有效flatMap,但在返回大量元素的情况下,它们都会按顺序处理。事实上,flatMap函数返回的流被强制进入顺序模式。请参阅ReferencePipeline.java的第 270 行。
“显而易见”要做的事情是使这个流并行,或者至少不强迫它是顺序的。这可能会也可能不会改善事情。它很可能会改善某些事情,但会使其他事情变得更糟。这里当然需要更好的政策,但我不确定它会是什么样子。
另请注意,通过向其提交运行管道的任务来强制并行流在您选择的 fork-join 池中运行的技术也是特定于实现的行为。它在 JDK 8 中以这种方式工作,但将来可能会改变。
| 归档时间: |
|
| 查看次数: |
3554 次 |
| 最近记录: |