【问题标题】:Generate infinite parallel stream生成无限并行流
【发布时间】:2017-08-23 06:09:35
【问题描述】:

问题

嗨,我有一个函数,我将返回无限的并行流(是的,在这种情况下它要快得多)生成的结果。所以很明显(或没有)我用过

Stream<Something> stream = Stream.generate(this::myGenerator).parallel()

它有效,但是......当我想限制结果时它不起作用(当流是连续的时一切都很好)。我的意思是,当我制作类似的东西时,它会产生结果

stream.peek(System.out::println).limit(2).collect(Collectors.toList())

但即使peek 输出产生超过 10 个元素,collect 仍然没有最终确定(生成速度很慢,所以这 10 个可能需要一分钟)......这是一个简单的例子。实际上,限制这些结果是一个未来,因为主要期望是在用户终止进程之前只获得比最近的结果更好的结果(其他情况是首先返回我可以通过抛出异常所做的事情,如果没有的话else 会有所帮助 [findFirst 没有,即使我在控制台上有更多元素并且在大约 30 秒内没有更多结果])。

所以,问题是……

如何复制?我的想法也是使用 RxJava,还有另一个问题 - 如何使用该工具(或其他)实现类似的结果。

代码示例

public Stream<Solution> generateSolutions() {
     final Solution initialSolution = initialSolutionMaker.findSolution();
     return Stream.concat(
          Stream.of(initialSolution),
          Stream.generate(continuousSolutionMaker::findSolution)
    ).parallel();
}

new Solver(instance).generateSolutions()
    .map(Solution::getPurpose)
    .peek(System.out::println)
    .limit(5).collect(Collectors.toList());

findSolution 的实现并不重要。 它有一些副作用,比如添加到解决方案 repo(singleton、sych 等),但仅此而已。

【问题讨论】:

标签: java java-8 rx-java java-stream


【解决方案1】:

正如already linked answer 中所解释的,高效并行流的关键点是使用已经具有固有大小的流源,而不是使用无大小甚至无限的流并在其上应用limit。注入大小在当前实现中根本不起作用,而确保已知大小不会丢失要容易得多。即使无法保留确切的尺寸,例如在应用filter 时,尺寸仍将作为估计尺寸。

所以不是

Stream.generate(this::myGenerator).parallel()
      .peek(System.out::println)
      .limit(2)
      .collect(Collectors.toList())

随便用

IntStream.range(0, /* limit */ 2).unordered().parallel()
         .mapToObj(unused -> this.myGenerator())
         .peek(System.out::println)
         .collect(Collectors.toList())

或者,更接近您的示例代码

public Stream<Solution> generateSolutions(int limit) {
    final Solution initialSolution = initialSolutionMaker.findSolution();
    return Stream.concat(
         Stream.of(initialSolution),
         IntStream.range(1, limit).unordered().parallel()
               .mapToObj(unused -> continuousSolutionMaker.findSolution())
   );
}

new Solver(instance).generateSolutions(5)
    .map(Solution::getPurpose)
    .peek(System.out::println)
    .collect(Collectors.toList());

【讨论】:

  • 你让这听起来很容易......我还没想过。
  • 但是还有一个简单的问题——限制是一个特性——它必须有相同的算法。顺便说一句,当我希望它受到限制时,我已经解决了这个问题,只需添加sequential().limit(1)
  • “特征”不需要与另一个特征具有相同的算法。 limit 在并行执行中表现不佳的事实已被记录,因此虽然无法在内部优化此特定场景可能确实令人惊讶,但它与文档一致。调用sequential().limit(1) 是巴洛克式的,您可以通过首先删除.parallel() 来实现相同的目的。尽管如此,即使对于某些操作(例如这对toArray() 产生了影响……
【解决方案2】:

不幸的是,这是预期的行为。我记得我至少看过两个关于这个问题的话题,这里是one of them

这个想法是Stream.generate 创建一个unordered infinite stream 并且limit 不会引入SIZED 标志。因此,当您在该 Stream 上生成 parallel 执行时,各个任务必须同步它们的执行以查看它们是否已达到该限制;到同步发生时,可能已经处理了多个元素。例如:

 Stream.iterate(0, x -> x + 1)
            .peek(System.out::println)
            .parallel()
            .limit(2)
            .collect(Collectors.toList());

还有这个:

IntStream.of(1, 2, 3, 4)
            .peek(System.out::println)
            .parallel()
            .limit(2)
            .boxed()
            .collect(Collectors.toList());

将始终在List (Collectors.toList) 中生成两个元素,并且将始终 也输出两个元素(通过peek)。

另一方面:

Stream<Integer> stream = Stream.generate(new Random()::nextInt).parallel();

List<Integer> list = stream
            .peek(x -> {
                System.out.println("Before " + x);
            })
            .map(x -> {
                System.out.println("Mapping x " + x);
                return x;
            })
            .peek(x -> {
                System.out.println("After " + x);
            })
            .limit(2)
            .collect(Collectors.toList());

将在List 中生成两个元素,但它可能会处理更多稍后将被limit 丢弃的元素。这就是您在示例中实际看到的内容。

唯一合理的方法(据我所知)是创建一个自定义拆分器。我没有写很多,但这是我的尝试:

 static class LimitingSpliterator<T> implements Spliterator<T> {

    private int limit;

    private final Supplier<T> generator;

    private LimitingSpliterator(Supplier<T> generator, int limit) {
        Preconditions.checkArgument(limit > 0);
        this.limit = limit;
        this.generator = Objects.requireNonNull(generator);
    }

    @Override
    public boolean tryAdvance(Consumer<? super T> consumer) {
        if (limit == 0) {
            return false;
        }
        T nextElement = generator.get();
        --limit;
        consumer.accept(nextElement);
        return true;
    }

    @Override
    public LimitingSpliterator<T> trySplit() {

        if (limit <= 1) {
            return null;
        }

        int half = limit >> 1;
        limit = limit - half;
        return new LimitingSpliterator<>(generator, half);
    }

    @Override
    public long estimateSize() {
        return limit >> 1;
    }

    @Override
    public int characteristics() {
        return SIZED;
    }
}

而用法是:

 StreamSupport.stream(new LimitingSpliterator<>(new Random()::nextInt, 7), true)
            .peek(System.out::println)
            .collect(Collectors.toList());

【讨论】:

  • 好的,谢谢您的回答。但是,这只是解释了为什么我的解决方案有错误。但是,正如我所提到的,限制只是一个 功能 - 它必须寻找无限的解决方案,所以我无法说明该解决方案。任何其他想法(例如反应性?)
  • 顺便说一句,奇怪的是,当我添加 limit(2) 时,它会返回大约 1000 个解决方案,例如,当我将它与另一个 limit(2) 链接时,它可能会返回 8-10(随机),但经过另一个链接没有区别......
  • @Azbesciak 1) 不 - 我不知道是否有办法使用 RxJava 和 2) 那些 limit.limit 可能只是简单的 samples 你得到:行为未定义,对此没有任何保证...
  • 这比编写自定义Spliterator 容易得多。请记住,映射操作对可拆分性或 SIZED 特性没有影响。所以规范的解决方案总是会导致IntStream.range(0, limit).map…(unused -&gt; actualGeneratorFunction)(或LongStream…,分别)......
猜你喜欢
  • 1970-01-01
  • 2019-06-11
  • 2021-01-27
  • 2021-04-11
  • 2013-12-13
  • 2020-05-17
  • 2016-05-13
  • 2023-03-14
相关资源
最近更新 更多