【问题标题】:Bug in parallelStream in javajava中parallelStream中的错误
【发布时间】:2021-05-31 08:43:04
【问题描述】:

谁能告诉我为什么会发生这种情况以及这是预期的行为还是错误

List<Integer> a = Arrays.asList(1,1,3,3);

a.parallelStream().filter(Objects::nonNull)
        .filter(value -> value > 2)
        .reduce(1,Integer::sum)

回答:10

但如果我们使用 stream 而不是 parallelStream 我得到了正确和预期 answer 7

【问题讨论】:

  • 检查documentation of reduce: "...标识值必须是累加器函数的标识。这意味着对于所有xaccumulator.apply(identity, x) 是相等的到x..." - 不是1Integer::sum
  • 为了能够使用并行流,减少是分部分完成的,每个部分都添加一个(可能每个元素一次) - 尝试使用0,结果应该是正确的跨度>

标签: java filter stream reduce parallelstream


【解决方案1】:

reduce 的第一个参数称为“identity”而不是“initialValue”。

1根据加法是没有身份的。 1 是乘法的恒等式。

如果你想对元素求和,你需要提供0


Java 使用“identity”而不是“initialValue”,因为这个小技巧可以轻松并行化 reduce


在并行执行中,每个线程将在流的一部分上运行reduce,当线程完成后,它们将使用完全相同的reduce函数进行组合。

虽然它看起来像这样:

mainThread:
  start thread1;
  start thread2;
  wait till both are finished;

thread1:
  return sum(1, 3); // your reduce function applied to a part of the stream

thread2:
  return sum(1, 3);

// when thread1 and thread2 are finished:
mainThread:
  return sum(sum(1, resultOfThread1), sum(1, resultOfThread2));
  = sum(sum(1, 4), sum(1, 4))
  = sum(5, 5)
  = 10

我希望你能看到,发生了什么以及为什么结果不是你所期望的。

【讨论】:

  • 但是为什么在流的情况下结果不同呢?
  • 因为非并行流按顺序工作。我会在几分钟后在我的回答中提供一个例子
  • 感谢@Benjamin,这很有帮助。
  • 有趣的是:算法有点不同 - 即使值被过滤掉,identity 也会被添加:IntStream.range(0, 100).parallel().filter(v -&gt; v &lt; 0).reduce(1,Integer::sum)36 使用 jshell 1.15)
  • @user15244370 是的。结果取决于正在使用的线程数。您的示例给出:4 用于 1 个线程,12 用于 2 个线程,20 用于 4 个线程,36 用于 8 个线程。我猜除了线程数之外还有一个块大小,所以每个线程一次只计算 N 个元素,完成后它会得到下一个 N 个元素。
猜你喜欢
  • 2014-07-20
  • 2018-11-24
  • 2014-05-27
  • 2013-11-01
  • 2023-03-26
  • 1970-01-01
  • 2021-04-06
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多