【问题标题】:Parallel Infinite Java Streams run out of Memory并行无限 Java 流内存不足
【发布时间】:2020-05-17 00:56:04
【问题描述】:

我试图理解为什么下面的 Java 程序给出了OutOfMemoryError,而没有.parallel() 的相应程序却没有。

System.out.println(Stream
    .iterate(1, i -> i+1)
    .parallel()
    .flatMap(n -> Stream.iterate(n, i -> i+n))
    .mapToInt(Integer::intValue)
    .limit(100_000_000)
    .sum()
);

我有两个问题:

  1. 这个程序的预期输出是什么?

    如果没有.parallel(),这似乎只是输出sum(1+2+3+...),这意味着它只是“卡在”flatMap 中的第一个流,这是有道理的。

    对于并行,我不知道是否有预期的行为,但我的猜测是它以某种方式交错了第一个 n 左右的流,其中 n 是并行工作者的数量。根据分块/缓冲行为,它也可能略有不同。

  2. 是什么导致它耗尽内存? 我特别想了解这些流是如何在后台实现的。

    我猜有什么东西阻塞了流,所以它永远不会完成并且能够摆脱生成的值,但我不太清楚事物的评估顺序和缓冲发生的位置。

编辑:如果相关,我使用的是 Java 11。

Editt 2: 显然,即使对于简单的程序 IntStream.iterate(1,i->i+1).limit(1000_000_000).parallel().sum(),也会发生同样的事情,所以它可能与 limit 而不是 flatMap 的懒惰有关。

【问题讨论】:

  • parallel() 内部使用 ForkJoinPool。我猜 ForkJoin 框架是从 Java 7 开始的 Java

标签: java java-stream out-of-memory lazy-evaluation


【解决方案1】:

您说“但我不太清楚事物的评估顺序以及缓冲发生的位置”,这正是并行流的意义所在。评估顺序未指定。

您的示例的一个关键方面是.limit(100_000_000)。这意味着实现不能只对任意值求和,而必须对 前 100,000,000 个数字求和。请注意,在参考实现中,.unordered().limit(100_000_000) 不会改变结果,这表明无序情况没有特殊实现,但这是一个实现细节。

现在,当工作线程处理元素时,它们不能只是总结它们,因为它们必须知道允许使用哪些元素,这取决于在其特定工作负载之前有多少元素。由于此流不知道大小,因此只有在处理了前缀元素时才能知道这一点,而无限流永远不会发生这种情况。因此,工作线程暂时保持缓冲,此信息可用。

原则上,当工作线程知道它处理了最左边的工作块时,它可以立即对元素求和,对它们进行计数,并在达到限制时发出结束信号。所以 Stream 可能会终止,但这取决于很多因素。

在您的情况下,一个合理的情况是其他工作线程分配缓冲区的速度比最左边的作业计算的速度快。在这种情况下,时间的细微变化可能会使流偶尔返回一个值。

当我们减慢除处理最左边块的工作线程之外的所有工作线程时,我们可以使流终止(至少在大多数运行中):

System.out.println(IntStream
    .iterate(1, i -> i+1)
    .parallel()
    .peek(i -> { if(i != 1) LockSupport.parkNanos(1_000_000_000); })
    .flatMap(n -> IntStream.iterate(n, i -> i+n))
    .limit(100_000_000)
    .sum()
);

¹我关注 a suggestion by Stuart Marks 在谈论遭遇顺序而不是处理顺序时使用从左到右的顺序。

【讨论】:

  • 非常好的答案!我想知道是否存在所有线程都开始运行 flatMap 操作的风险,并且没有分配到实际清空缓冲区(求和)?在我的实际用例中,无限流是文件太大而无法保存在内存中。我想知道如何重写流以降低内存使用率?
  • 你在使用Files.lines(…)吗?它在 Java 9 中得到了显着改进。
  • 它似乎只是调用了BufferedReader.lines,这只是一个简单行迭代器的拆分器包装器。所以我希望它会继续进入我的记忆,直到我用完空间,就像iterate一样。但我当然必须对其进行测试。
  • 这就是它在 Java 8 中所做的事情。在较新的 JRE 中,在某些情况下(不是默认文件系统、特殊字符集或大于 @987654329 的大小),它仍会回退到 BufferedReader.lines() @)。如果其中之一适用,自定义解决方案可能会有所帮助。这将值得一个新的问答......
  • 什么是外流,文件流?它有可预测的大小吗?
【解决方案2】:

我最好的猜测是添加 parallel() 会改变 flatMap() 的内部行为 already had problems being evaluated lazily before

[JDK-8202307] Getting a java.lang.OutOfMemoryError: Java heap space when calling Stream.iterator().next() on a stream which uses an infinite/very big Stream in flatMap 报告了您收到的 OutOfMemoryError 错误。如果您查看票证,它或多或少与您获得的堆栈跟踪相同。工单因无法修复而关闭,原因如下:

iterator()spliterator() 方法是在无法使用其他操作时使用的“逃生舱口”。它们有一些限制,因为它们将流实现的推模型转变为拉模型。这种过渡在某些情况下需要缓冲,例如当一个元素(平面)映射到两个或多个元素时。这将使流实现显着复杂化,可能以牺牲常见情况为代价,以支持背压的概念,以传达有多少元素要通过元素生产的嵌套层。

【讨论】:

  • 这很有趣!推/拉转换需要缓冲,这可能会耗尽内存,这是有道理的。但是,在我的情况下,似乎只使用 push 应该可以正常工作,并且在剩余元素出现时简单地丢弃它们?或者你是说 flapmap 会导致创建一个迭代器?
【解决方案3】:

OOME 不是是由于流是无限的,而是由于 它不是的事实。

也就是说,如果你注释掉 .limit(...),它永远不会耗尽内存——当然,它也永远不会结束。

一旦拆分,流只能跟踪每个线程中累积的元素数量(看起来实际的累加器是Spliterators$ArraySpliterator#array)。

看起来你可以在没有flatMap 的情况下重现它,只需使用-Xmx128m 运行以下命令:

    System.out.println(Stream
            .iterate(1, i -> i + 1)
            .parallel()
      //    .flatMap(n -> Stream.iterate(n, i -> i+n))
            .mapToInt(Integer::intValue)
            .limit(100_000_000)
            .sum()
    );

但是,在注释掉 limit() 之后,它应该可以正常运行,直到您决定不使用笔记本电脑。

除了实际的实现细节,我认为是这样的:

使用limitsum reducer 希望前 X 个元素求和,因此没有线程可以发出部分和。每个“切片”(线程)都需要累积元素并传递它们。没有限制,就没有这样的约束,所以每个“切片”只会(永远)计算它得到的元素的部分总和,假设它最终会发出结果。

【讨论】:

  • “一旦拆分”是什么意思?限制会以某种方式拆分它吗?
  • @ThomasAhle parallel() 将在内部使用ForkJoinPool 来实现并行性。 Spliterator 将用于为每个ForkJoin 任务分配工作,我想我们可以将这里的工作单元称为“拆分”。
  • 但为什么只有在limit的情况下才会出现这种情况?
  • 奇怪的是,IntStream.iterate(1,i->i+1).parallel().sum() 运行良好(不会停止,但也不会崩溃);版本IntStream.iterate(1,i->i+1).limit(1000_000_000).parallel().sum() 内存不足。这表明 limit 正在做一些奇怪的事情,而不仅仅是对元素进行排序。
  • @ThomasAhle 在Integer.sum() 中设置断点,由IntStream.sum 减速器使用。您会看到无限制版本一直调用该函数,而限制版本永远不会在 OOM 之前调用它。
猜你喜欢
  • 2016-08-13
  • 2015-12-02
  • 1970-01-01
  • 1970-01-01
  • 2016-10-10
  • 2019-05-13
  • 1970-01-01
  • 1970-01-01
  • 2016-05-13
相关资源
最近更新 更多