【问题标题】:Will inner parallel streams be processed fully in parallel before considering parallelizing outer stream?在考虑并行化外部流之前,内部并行流是否会完全并行处理?
【发布时间】:2018-01-16 04:05:54
【问题描述】:

从这个link,我只是部分理解,至少在某些时候,Java 嵌套并行流存在问题。但是,我无法推断出以下问题的答案:

假设我有一个外部 srtream 和一个内部流,它们都使用并行流。事实证明,根据我的计算,如果内部流首先完全并行完成,然后(如果且仅cpu核心可用)做外部流。我认为这对于大多数人的情况都是正确的。所以我的问题是:

Java 会先并行执行内部流,然后再处理外部流吗?如果是这样,它是在编译时还是在运行时做出决定?如果在运行时,JIT 是否足够聪明地意识到如果内部流确实有超过足够的元素(例如数百个)而不是核心数(32),那么它肯定应该使用所有 32 个核心来处理在从外部流移动到下一个元素之前的内部流;但是,如果元素的数量很少(例如

【问题讨论】:

  • 你能举例说明你的意思吗?喜欢flatMap.parallel?streamA.... map(streamB.parallel...)
  • 流并行性几乎完全不可配置,它旨在自动运行。如果您需要优化并行性,我建议您根本不要使用流。无论如何,他们有很多开销。

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


【解决方案1】:

也许下面的示例程序可以解决这个问题:

IntStream.range(0, 10).parallel().mapToObj(i -> "outer "+i)
         .map(outer -> outer+"\t"+IntStream.range(0, 10).parallel()
            .mapToObj(inner -> Thread.currentThread())
            .distinct() // using the identity of the threads
            .map(Thread::getName) // just to be paranoid, as names might not be unique
            .sorted()
            .collect(Collectors.toList()) )
         .collect(Collectors.toList())
         .forEach(System.out::println);

当然,结果会有所不同,但我机器上的输出看起来是这样的:

outer 0 [ForkJoinPool.commonPool-worker-6]
outer 1 [ForkJoinPool.commonPool-worker-3]
outer 2 [ForkJoinPool.commonPool-worker-1]
outer 3 [ForkJoinPool.commonPool-worker-1, ForkJoinPool.commonPool-worker-4, ForkJoinPool.commonPool-worker-5]
outer 4 [ForkJoinPool.commonPool-worker-5]
outer 5 [ForkJoinPool.commonPool-worker-2, ForkJoinPool.commonPool-worker-4, ForkJoinPool.commonPool-worker-7, main]
outer 6 [main]
outer 7 [ForkJoinPool.commonPool-worker-4]
outer 8 [ForkJoinPool.commonPool-worker-2]
outer 9 [ForkJoinPool.commonPool-worker-7]

我们在这里可以看到,对于我的机器,有八个内核,七个工作线程正在参与工作,以利用所有内核,至于common pool,调用者线程也将参与工作,而不是仅仅等待完成。您可以清楚地看到输出中的main 线程。

此外,您可以看到外部流获得了完全的并行性,而一些内部流完全由单个线程处理。每个工作线程都至少为外部流的一个元素做出贡献。如果您将外部流的大小减少到核心数,您很可能会看到一个工作线程处理一个外部流元素,这意味着所有内部流完全按顺序执行。

但我使用了一个与内核数不匹配的数字,甚至不是它的倍数,来演示另一种行为。由于外部流处理的工作负载不均匀,即一些线程只处理一项,而另一些线​​程处理两项,这些空闲的工作线程执行工作窃取,贡献了剩余外部元素的内部流处理。

这种行为背后有一个简单的理由。当外部流的处理开始时,它并不知道它将是一个“外部流”。它只是一个并行流,除了处理它之外,没有办法找出这是否是外部流,直到其中一个函数开始另一个流操作。但是将并行处理推迟到可能永远不会到来的这一点是没有意义的。

除此之外,我强烈反对您的假设“如果内部流首先完全并行完成,性能会更高 [...]”。对于典型的用例,我宁愿反过来期待它,阅读,期待这样做的优势,就像它已经实现的那样。但是,如前一段所述,无论如何都没有合理的方法来实现并行处理内部流的偏好。

【讨论】:

  • 我认为首先并行处理内部流会更高效,因为我的内部流通常是集合,其中集合类中有一些“大”哈希图正好适合 CPU三级缓存。因此,如果内部流都并行运行,那么这个“大”哈希映射将适合 L3 缓存,并且所有内核都可以访问它(并行)。但是,如果内部流按顺序运行(因为外部流是并行化的),那么每个内部流线程将竞争缓存大小的 1/32(在 32 核机器上)。想法?
  • 想法@Holger?
  • A HashMap 不是内存块。 HashMap,它的后备数组,每个Entry 实例,引用的键对象和值对象都是不同的对象,不一定在内存中相邻,甚至不能保证在同一个区域。尝试一次将它们全部填充到 L3 缓存中没有任何优势。你最终会导致每个线程仍然使用 L3 缓存的 ¹/₃₂,无论它属于同一个 HashMap 还是不同的 HashMaps,因为并行处理的基本原则是每个线程处理不同的数据块。
【解决方案2】:

根据我刚才写的小测试答案是no(约Would Java execute inner stream all in parallel first, and then work on outerstream)。请注意,默认情况下,在我的机器上有 4 个线程用于流操作。

    List<Integer> first = List.of(1, 2, 3, 4);
    List<Integer> second = List.of(5, 6, 7, 8);

    first.stream().parallel()
            .peek(x -> {
                System.out.println("first : " + x + " " + Thread.currentThread().getName());
            })
            .map(x -> second.stream().parallel().peek(y -> {

                System.out.println("second : " + y + " " + Thread.currentThread().getName());

            }).collect(Collectors.toList()))
            .filter(x -> true)
            .collect(Collectors.toList());

从输出中可以看出,内部流没有先执行。您可以增加每个流中的元素数量以获得更准确的输出(交错“第一”和“第二” - 不知道这是否正确)。

但是这里还有一些让我印象深刻的东西......上面的例子如何不阻塞是我无法理解的。只有 4 个线程和 4 个元素,所有线程都在等待内部流处理;但是ForkJoinPool 没有可用的线程可供使用 - 那么它是如何工作的呢? 您提供的链接(@Holger 的回答)表明创建的线程将比您实际请求的更多线程。但是他们的名字从输出中消失了......

【讨论】:

  • 看来,你误会了。 F/J 将创建补偿线程,如果线程被故意阻塞。但是,在 Stream API 的情况下,调用者线程不会被阻塞,而是会参与处理。早期的 Java 8 版本似乎存在这种工作窃取的问题,但这似乎已经得到了改进。在工作线程中等待CompletableFuture 时,您仍然可以看到正在运行的补偿线程。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多