【问题标题】:Why does stream parallel() not use all available threads?为什么流并行()不使用所有可用线程?
【发布时间】:2020-01-21 12:38:37
【问题描述】:

我尝试使用 Java8(1.8.0_172) stream.parallel() 并行运行 100 个 Sleep 任务,该任务在具有 100 多个可用线程的自定义 ForkJoinPool 中提交。每个任务会休眠 1s。我预计整个工作将在大约 1 秒后完成,因为 100 个睡眠可以并行完成。但是我观察到 7 秒的运行时间。

    @Test
    public void testParallelStream() throws Exception {
        final int REQUESTS = 100;
        ForkJoinPool forkJoinPool = null;
        try {
            // new ForkJoinPool(256): same results for all tried values of REQUESTS
            forkJoinPool = new ForkJoinPool(REQUESTS);
            forkJoinPool.submit(() -> {

                IntStream stream = IntStream.range(0, REQUESTS);
                final List<String> result = stream.parallel().mapToObj(i -> {
                    try {
                        System.out.println("request " + i);
                        Thread.sleep(1000);
                        return Integer.toString(i);
                    } catch (InterruptedException e) {
                        throw new RuntimeException(e);
                    }
                }).collect(Collectors.toList());
                // assertThat(result).hasSize(REQUESTS);
            }).join();
        } finally {
            if (forkJoinPool != null) {
                forkJoinPool.shutdown();
            }
        }
    }

输出指示在 1 秒的暂停之前执行 ~16 个流元素,然后再执行 ~16 个,依此类推。所以看起来即使 forkjoinpool 是用 100 个线程创建的,也只有大约 16 个线程被使用。

当我使用超过 23 个线程时,这种模式就会出现:

1-23 threads: ~1s
24-35 threads: ~2s
36-48 threads: ~3s
...
System.out.println(Runtime.getRuntime().availableProcessors());
// Output: 4

【问题讨论】:

  • 你的Runtime.getRuntime().availableProcessors() 输出是什么?
  • availableProcessors() == 4,我添加到描述中
  • 顺序执行需要多长时间?
  • 这个问题可能观察到同样的事情(没有答案)stackoverflow.com/questions/49068119
  • 这种使并行流使用不同线程池的技巧是一种未记录的实现副作用,并且不打算以这种方式工作。所以实现并不关心不同并行的可能性。

标签: java multithreading java-8 java-stream forkjoinpool


【解决方案1】:

由于 Stream 实现对 Fork/Join 池的使用是一个实现细节,因此强制它使用不同的 Fork/Join 池的技巧也没有记录,并且似乎是偶然的,即有一个 hardcoded constant 确定实际的并行度,取决于默认池的并行度。因此,最初并没有预见到使用不同的池。

但是,已经认识到使用具有不适当目标并行度的不同池是一个错误,即使没有记录此技巧,请参阅JDK-8190974

它已在 Java 10 中修复并向后移植到 Java 8,更新 222。

所以一个简单的解决方案世界正在更新 Java 版本。

您还可以更改默认池的并行度,例如

System.setProperty("java.util.concurrent.ForkJoinPool.common.parallelism", "100");

在进行任何 Fork/Join 活动之前。

但这可能会对其他并行操作产生意想不到的影响。

【讨论】:

    【解决方案2】:

    在您编写它时,您让流决定执行的并行度。

    您的效果是ArrayList.parallelStream 试图通过平均拆分数据来超越您,而不考虑可用线程的数量。这对于 CPU-Bound 操作很有用,在这种操作中,线程数多于 CPU 内核数并没有什么用处,但不适用于需要等待 IO 的进程。

    为什么不按顺序将所有项目强制提供给 ForkJoinPool,从而强制使用所有可用线程?

            IntStream stream = IntStream.range(0, REQUESTS);
            List<ForkJoinTask<String>> results
                    = stream.mapToObj(i -> forkJoinPool.submit(() -> {
    
                try {
                    System.out.println("request " + i);
                    Thread.sleep(1000);
                    return Integer.toString(i);
                } catch (InterruptedException e) {
                    throw new RuntimeException(e);
                }
            })).collect(Collectors.toList());
            results.forEach(ForkJoinTask::join);
    

    这在我的机器上不到两秒钟。

    【讨论】:

    • 重点是了解并行流的工作原理,以及编写代码时要注意什么。所以你的解决方案没有错,但问题是专门关于使用并行流的解决方案不起作用。
    猜你喜欢
    • 2016-08-25
    • 1970-01-01
    • 2015-06-07
    • 1970-01-01
    • 2018-03-19
    • 1970-01-01
    • 1970-01-01
    • 2017-10-16
    相关资源
    最近更新 更多