【问题标题】:Is there a way to force parallelStream() to go parallel?有没有办法强制 parallelStream() 并行?
【发布时间】:2017-12-01 16:09:06
【问题描述】:

如果输入大小太小,库automatically serializes the execution of the maps in the stream,但这种自动化不会也无法考虑映射操作的繁重程度。有没有办法强制 parallelStream() 实际并行化 CPU heavy 映射?

【问题讨论】:

  • 您的链接问题已经包含您在 cmets 中问题的答案(由尊敬的 Brian Goetz 提供)。
  • 好吧,正如解释的那样,你不能强迫它。而是使用执行器。您添加冗余元素的解决方法是一个非常可怕的黑客攻击。
  • 我明白这一点,但有时您需要使用正确的工具来完成正确的工作,而不是仅仅因为您认为必须使用错误的工具而坚持使用错误的工具并四处乱窜。如果你只有一把锤子,那么一切看起来都像钉子,听起来你的锤子就是流 api。
  • 我不是在谈论 Stream API 的可能限制。我说的是您选择次优解决方案的唯一原因是您想使用 Stream API。对于软件开发人员来说,这不是一个很好的品质。
  • 好吧,so far 看来你在 JDK9 中也不会get what you want

标签: java concurrency parallel-processing java-8 java-stream


【解决方案1】:

似乎存在根本性的误解。链接的问答讨论了流显然不能并行工作,因为 OP 没有看到预期的加速。结论是,如果工作负载太小,并行处理没有任何好处没有自动回退到顺序执行。

实际上恰恰相反。如果您请求并行,您将获得并行,即使它实际上会降低性能。在这种情况下,实现不会切换到可能更有效的顺序执行。

因此,如果您确信每个元素的工作负载足够高,足以证明使用并行执行的合理性(无论元素数量如何),您都可以简单地请求并行执行。

很容易证明:

Stream.of(1, 2).parallel()
      .peek(x -> System.out.println("processing "+x+" in "+Thread.currentThread()))
      .forEach(System.out::println);

On Ideone,它会打印出来

processing 2 in Thread[main,5,main]
2
processing 1 in Thread[ForkJoinPool.commonPool-worker-1,5,main]
1

但消息和详细信息的顺序可能会有所不同。甚至有可能在某些环境中,两个任务都可能碰巧由同一​​个线程执行,如果它可以在另一个线程开始接收它之前巩固第二个任务。但是,当然,如果任务足够昂贵,就不会发生这种情况。重要的一点是,整个工作负载已被拆分并排入队列,可能会被其他工作线程处理。

如果上面的简单示例在您的环境中发生单线程执行,您可以像这样插入模拟工作负载:

Stream.of(1, 2).parallel()
      .peek(x -> System.out.println("processing "+x+" in "+Thread.currentThread()))
      .map(x -> {
           LockSupport.parkNanos("simulated workload", TimeUnit.SECONDS.toNanos(3));
           return x;
        })
      .forEach(System.out::println);

然后,如果“每个元素的处理时间”已经足够长了。


更新:误解可能是由 Brian Goetz 的误导性陈述引起的:“在您的情况下,您的输入集太小而无法分解”。

必须强调的是,这不是Stream API的通用属性,而是已经使用的MapHashMap 有一个后备数组,条目根据其哈希码分布在该数组中。将数组拆分为 n 个范围可能不会导致包含元素的平衡拆分,尤其是在只有两个的情况下。 HashMapSpliterator 的实现者认为在数组中搜索元素以获得完美平衡的拆分过于昂贵,并不是说拆分两个元素不值得。

由于HashMap 的默认容量是16,并且示例只有两个元素,我们可以说地图过大。简单地修复它也可以修复示例:

long start = System.nanoTime();

Map<String, Supplier<String>> input = new HashMap<>(2);
input.put("1", () -> {
    System.out.println(Thread.currentThread());
    LockSupport.parkNanos("simulated workload", TimeUnit.SECONDS.toNanos(2));
    return "a";
});
input.put("2", () -> {
    System.out.println(Thread.currentThread());
    LockSupport.parkNanos("simulated workload", TimeUnit.SECONDS.toNanos(2));
    return "b";
});
Map<String, String> results = input.keySet()
        .parallelStream().collect(Collectors.toConcurrentMap(
    key -> key,
    key -> input.get(key).get()));

System.out.println("Time: " + TimeUnit.NANOSECONDS.toMillis(System.nanoTime()- start));

在我的机器上打印

Thread[main,5,main]
Thread[ForkJoinPool.commonPool-worker-1,5,main]
Time: 2058

结论是 Stream 实现总是尝试使用并行执行,如果你请求它,不管输入大小。但这取决于输入的结构如何将工作负载分配给工作线程。事情可能会更糟,例如如果您从文件中流式传输行。

如果您认为平衡拆分的好处值得复制步骤的成本,您也可以使用new ArrayList&lt;&gt;(input.keySet()).parallelStream() 而不是input.keySet().parallelStream(),因为ArrayList 内的元素分布始终允许完美平衡拆分.

【讨论】:

  • @maverickabhi 您不能同时处理“按顺序”和“使用并行处理”。这些是相互矛盾的术语。当您想要的只是将数组写入文件时,并行处理根本没有任何好处。如果您有一个计算昂贵的中间操作的流,您可以尝试使用并行流和链forEachOrdered 作为终端操作以并行写入目标文件。但是根据实际操作,按顺序编写最终结果的成本可能仍然超过并行处理的任何好处。
  • @Kamel 你希望在代码中找到什么样的信息?
  • @Just 我不知道任何现实生活中parallelStream() 的行为与stream().parallel() 不同的例子。数组既没有 stream() 也没有 parallelStream()
  • @只需发布minimal reproducible example,而不是包含大量非标准方法和类的 sn-p。有在线测试器,您可以在其中运行代码来演示行为。也许,您想打开一个新问题,而不是用 cmets 淹没这个问题。
  • @Just parallelStream()stream().parallel() 之间没有区别,除了互联网上的单个用户在没有任何证据的情况下声称不这样做。要么证明你的主张,要么停止讨论。
猜你喜欢
  • 1970-01-01
  • 2011-03-28
  • 1970-01-01
  • 2014-04-17
  • 1970-01-01
  • 2018-02-25
  • 2022-10-06
  • 2021-08-07
  • 2014-01-02
相关资源
最近更新 更多