【发布时间】: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 等),但仅此而已。
【问题讨论】:
-
您能否提供Minimal, Verifiable, Complete example 来澄清您的问题?
-
添加了示例代码。
标签: java java-8 rx-java java-stream