【问题标题】:parallel processing with infinite stream in JavaJava中无限流的并行处理
【发布时间】:2016-05-13 09:04:42
【问题描述】:

为什么下面的代码不打印任何输出,而如果我们删除并行,它会打印 0、1?

IntStream.iterate(0, i -> ( i + 1 ) % 2)
         .parallel()
         .distinct()
         .limit(10)
         .forEach(System.out::println);

虽然我知道理想情况下应该将 limit 放在 distinct 之前,但我的问题更多与添加并行处理造成的差异有关。

【问题讨论】:

  • 在我的机器上,这会以 100% 的 CPU 使用率锁定 3/4 个内核,而不会产生答案!我认为这可能是limit和parallel交互的一个bug。
  • @Straw1239 - 同意这应该被视为一个错误;这段代码没有立即打印出一些东西,没有明显的原因。
  • 它尝试在distinct op 中缓冲整个流内容。即使没有limit,也会发生这种情况。虽然很明显,处理后续块的线程必须等待它们的前一个流来获得有序流,但第一个不需要等待......
  • @bayou.io:这是一个已知现象,没有看到 OOME,但有一个进程挂在仅执行 GC 活动中。只需监控 JVM 活动、已用内存与最大堆之间的比率以及 CPU 使用率与 GC 活动之间的比率。不过,当 GC 开始发疯时,大多数监控方法也会停止工作……
  • @Straw1239:bayou.io 谈到使用i->i+1 函数而不是i->(i+1)%2。当不使用%2 时,缓冲区真的开始耗尽所有内存。

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


【解决方案1】:

真正的原因是 ordered parallel .distinct() 是文档中described 的完整屏障操作:

在并行管道中保持distinct() 的稳定性相对昂贵(要求操作充当完整的屏障,具有大量缓冲开销),并且通常不需要稳定性。

“全屏障操作”是指必须先执行所有上游操作,然后才能启动下游。 Stream API 中只有两个完整的屏障操作:.sorted()(每次)和.distinct()(在有序并行情况下)。由于您向.distinct() 提供了非短路无限流,因此您最终会出现无限循环。根据合同.distinct() 不能只是以任何顺序向下游发出元素:它应该始终发出第一个重复元素。虽然理论上可以更好地实现并行有序.distinct(),但实现起来会复杂得多。

至于解决方案,@user140547 是对的:在.distinct() 之前添加.unordered() 这会将distinct() 算法切换为无序算法(它只使用共享ConcurrentHashMap 来存储所有观察到的元素并将每个新元素发送到下游)。请注意,添加.unordered() .distinct() 将无济于事。

【讨论】:

    【解决方案2】:

    我知道代码不正确,并且正如解决方案中所建议的那样,如果我们在 distinct 之前移动限制,我们将不会有无限循环。

    并行函数是使用fork和join概念来分配工作,它为工作分配所有可用的线程,而不是单个线程。

    我们正确地期待无限循环,因为多个线程无限地处理数据并且没有任何东西阻止它们,因为 10 的限制永远不会在不同之后达到。

    它可能会一直尝试分叉并且从不尝试加入以使其前进。但我仍然认为它是 java 中的一个缺陷,最重要的是。

    【讨论】:

    • limit 移到distinct 之前会改变语义。除此之外,您可以使用limit(2) 使您的代码正确,但问题仍然存在。
    【解决方案3】:

    这段代码有一个大问题,即使没有并行: 在 .distinct() 之后,流将只有 2 个元素——因此限制永远不会生效——它将打印这两个元素,然后无限期地继续浪费你的 CPU 时间。不过,这可能是您的本意。

    使用并行和限制,我认为由于工作的划分方式,问题会更加严重。并行流代码我还没有一路追踪,但这是我的猜测:

    并行代码在多个线程之间划分工作,所有线程都无限期地运行,因为它们永远不会填满它们的配额。系统可能会等待每个线程完成,因此它可以组合它们的结果以确保按顺序区分 - 但在您提供的情况下永远不会发生这种情况。

    在没有顺序要求的情况下,每个工作线程的结果在与全局独特性集进行检查后可以立即使用。

    没有限制,我怀疑使用不同的代码来处理无限流:而不是等待所需的 10 填满,结果报告为已发现。它有点像制作一个报告 hasNext() = true 的迭代器,首先产生 0,然后是 1,然后 next() 调用永远挂起而不产生结果 - 在并行情况下,某些东西正在等待多个报告,因此它可以正确组合/在输出之前订购它们,而在串行情况下,它会执行它可以挂起的操作。

    我将尝试找出有无 distinct() 或 limit() 的调用堆栈的确切差异,但到目前为止,在相当复杂的流库调用序列中导航似乎非常困难。

    【讨论】:

    • 您一定尝试过不同的实现。对于 Oracle 的 1.8.0_601.8.0_65,无论您使用 .limit(10) 还是 .limit(2) 都没有关系,即使省略 limit 也不会改变行为。
    【解决方案4】:

    Stream.iterate 返回“无限顺序有序流”。因此,使顺序流并行并没有太大用处。

    根据Stream package的描述:

    对于并行流,放宽排序约束有时可以提高执行效率。如果元素的顺序不相关,则可以更有效地实现某些聚合操作,例如过滤重复项 (distinct()) 或分组缩减 (Collectors.groupingBy())。类似地,本质上与遇到顺序相关的操作,例如 limit(),可能需要缓冲以确保正确排序,从而破坏并行性的好处。在流具有遇到顺序但用户并不特别关心该遇到顺序的情况下,使用 unordered() 显式地对流进行降序可能会提高某些有状态或终端操作的并行性能。然而,大多数流管道,例如上面的“块权重之和”示例,即使在排序约束下仍然可以高效地并行化。

    这似乎是你的情况,使用 unordered(),它打印 0,1。

        IntStream.iterate(0, i -> (i + 1) % 2)
                .parallel()
                .unordered()
                .distinct()
                .limit(10)
                .forEach(System.out::println);
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-05-17
      • 1970-01-01
      • 2019-09-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多