【问题标题】:How to lazily evaluate nested flatMap如何懒惰地评估嵌套的 flatMap
【发布时间】:2021-08-12 10:49:03
【问题描述】:

我试图从两个可能无限的流中变出一个笛卡尔积,然后我通过limit() 对其进行限制。

到目前为止,这(大约)是我的策略:

@Test
void flatMapIsLazy() {
        Stream.of("a", "b", "c")
            .flatMap(s -> Stream.of("x", "y")
                .flatMap(sd -> IntStream.rangeClosed(0, Integer.MAX_VALUE)
                    .mapToObj(sd::repeat)))
            .map(s -> s + "u")
            .limit(20)
            .forEach(System.out::println);
}

这不起作用。

显然,我的第二个流在第一次在管道上使用时会在现场进行最终评估。它不会产生我可以按照自己的节奏消费的惰性流。

我认为ReferencePipeline#flatMap 的这段代码中的.forEach 是罪魁祸首:

@Override
public void accept(P_OUT u) {
    try (Stream<? extends R> result = mapper.apply(u)) {
        if (result != null) {
            if (!cancellationRequestedCalled) {
               result.sequential().forEach(downstream);
            }
            else {
                var s = result.sequential().spliterator();
                do { } while (!downstream.cancellationRequested() && s.tryAdvance(downstream));
            }
        }
    }
}

我希望上面的代码返回 20 个元素,如下所示:

a
ax
axx
axxx
axxxx
...
axxxxxxxxxxxxxxxxxxx

但它却因OutOfMemoryError 而崩溃,因为嵌套的flatMap 的很长的Stream 被急切地评估(??)并用重复字符串的不必要副本填满了我的记忆。如果提供了值 3 而不是 Integer.MAX_VALUE,将相同的限制保持在 20,则预期输出将改为:

a
ax
axx
axxx
a
ay
ayy
ayyy
b
bx
bxx
bxxx
...
(up until 20 lines)

编辑:此时我刚刚使用惰性迭代器推出了自己的实现。不过,我认为应该有一种方法可以使用纯 Streams 来做到这一点。

编辑 2:这已被承认为 Java https://bugs.java.com/bugdatabase/view_bug.do?bug_id=JDK-8267758%20 中的错误票

【问题讨论】:

  • 流大小不谈,你试过只运行一次"x".repeat(Integer.MAX_VALUE) 吗?在我的机器上,我得到了一个 OOM。也许这只是你在这里的一个坏例子,但你不能指望它会起作用。
  • 除此之外,.flatMap(s -&gt; second) 不能工作。您正在尝试重用流。这几乎肯定会给你一个非法状态异常。
  • 对原始查询更真实的代码版本可以是:Stream.of("a", "b", "c").flatMap(s -&gt; Stream.of("x", "y").flatMap(sd -&gt; IntStream.rangeClosed(0, Integer.MAX_VALUE).mapToObj(sd::repeat))).map(s -&gt; s + "u").limit(20).forEach(System.out::println);,这会导致 OOM。请注意,它有嵌套的 flatMap 调用。
  • @ernest_k 是的,就是这样。我更改了问题代码!谢谢! :)
  • 我假设您将提供一些值(如限制方法的整数)来控制输出的大小。查看多个此类值的预期输出会很有用。

标签: java java-stream lazy-evaluation cartesian-product


【解决方案1】:

正如您已经写的,这已被接受为一个错误。也许,它会在未来的 Java 版本中得到解决。

但即使是现在也可能有解决方案。它不是很优雅,只有在外部流中的元素数量和限制足够小时才可能可行。但它会在这些限制下工作。

让我先稍微修改一下您的示例,将外部 flatMap 转换为两个操作(一个 map 和一个带有标识的 flatMap,只做展平):

Stream.of("a", "b", "c")
      .map(s -> Stream.of("x", "y")
            .flatMap(sd -> IntStream.rangeClosed(0, Integer.MAX_VALUE)
                  .mapToObj(sd::repeat)))
      .flatMap(s -> s)
      .map(s -> s + "u")
      .limit(20)
      .forEach(System.out::println);

我们可以很容易地看到,每个内部流中我们需要的元素不超过 20 个。所以我们可以将每个流限制为这个数量的元素。这将起作用(您应该使用变量或常量作为限制):

Stream.of("a", "b", "c")
      .map(s -> Stream.of("x", "y")
            .flatMap(sd -> IntStream.rangeClosed(0, Integer.MAX_VALUE)
                  .mapToObj(sd::repeat)))
      .flatMap(s -> s.limit(20))            // limit each inner stream
      .map(s -> s + "u")
      .limit(20)
      .forEach(System.out::println);

当然这样还是会产生过多的中间结果,但在上述限制下可能问题不大。

【讨论】:

  • 太好了!在我的真实测试代码中它不起作用,因为我flatMap 和limit 在非常不同的位置和堆栈深度,因为内部和外部流是在不同的类中生成的。可能我可以通过调用使限制上下移动,但有点违背了目的。无论如何,感谢您的时间和回答。
猜你喜欢
  • 2017-04-30
  • 1970-01-01
  • 2019-01-07
  • 2017-01-13
  • 2021-07-10
  • 1970-01-01
  • 2016-05-12
  • 2014-10-10
  • 1970-01-01
相关资源
最近更新 更多