【发布时间】:2015-07-13 02:53:29
【问题描述】:
我有一个一般的 Streams API 问题,我想“有效地”解决。假设我有一个(可能非常大,可能是无限的)流。我想以某种方式对其进行预处理,例如,过滤掉一些项目,并对一些项目进行变异。假设这种预处理是复杂的、时间和计算密集型的,所以我不想做两次。
接下来,我想对项目序列执行两组不同的操作,并使用不同的流类型构造处理每个不同序列的远端。对于无限流,这将是一个 forEach,对于一个有限流,它可能是一个收集器或其他什么。
显然,我可能会将中间结果收集到一个列表中,然后从该列表中拖出两个单独的流,分别处理每个流。这适用于有限流,但 a) 它看起来“丑陋”,b) 对于非常大的流来说可能不切实际,而且完全不适用于无限流。
我想我可以将 peek 用作一种“三通”。然后,我可以对 peek 下游的结果执行一系列处理,并以某种方式强制消费者在 peek 中执行“其他”工作,但现在第二条路径不再是流。
我发现我可以创建一个 BlockingQueue,使用 peek 将内容推送到该队列中,然后从队列中获取一个流。这似乎是一个好主意,实际上效果很好,尽管我不明白流是如何关闭的(它实际上是这样,但我看不到如何关闭)。这是说明这一点的示例代码:
List<Student> ls = Arrays.asList(
new Student("Fred", 2.3F)
// more students (and Student definition) elided ...
);
BlockingQueue<Student> pipe = new LinkedBlockingQueue<>();
ls.stream()
.peek(s -> {
try {
pipe.put(s);
} catch (InterruptedException ioe) {
ioe.printStackTrace();
}
})
.forEach(System.out::println);
new Thread(
new Runnable() {
public void run() {
Map<String, Double> map =
pipe.stream()
.collect(Collectors.groupingBy(s->s.getName(),
Collectors.averagingDouble(s->s.getGpa())));
map.forEach(
(k,v)->
System.out.println(
"Students called " + k
+ " average " + v));
}
}).start();
那么,第一个问题是:有没有“更好”的方法来做到这一点?
第二个问题,BlockingQueue 上的流到底是如何关闭的?
干杯, 托比
【问题讨论】:
-
我认为,您所有的流都没有关闭。当流关闭时,将调用
onClose()处理程序。注册这样的处理程序并检查它是否被调用。 Stream 不会自动关闭:您应该使用 try-with-resources 语句或手动关闭它。可能您所说的“流已关闭”是指不同的意思?
标签: java-8 java-stream