没有isFinite()和isOrdered()方法,如何安全地使用Java Streams?

tkr*_*use 18 java java-stream

关于java方法是否应该返回Collections或Streams,存在一个问题,在该问题中Brian Goetz回答说,即使对于有限序列,通常也应该首选Streams。

但是在我看来,当前无法对来自其他地方的Streams进行许多操作,并且无法提供防御性的代码保护,因为Streams不会显示它们是无限的还是无序的。

如果并行是我要在Stream()上执行的操作的问题,我可以调用isParallel()进行检查或顺序执行,以确保计算是并行的(如果我记得的话)。

但是,如果有序性或无穷(大小)与我的程序的安全性有关,则我无法编写保护措施。

假设我使用一个实现此虚拟接口的库:

public interface CoordinateServer {
    public Stream<Integer> coordinates();
    // example implementations:
    // IntStream.range(0, 100).boxed()   // finite, ordered, sequential
    // final AtomicInteger atomic = new AtomicInteger();
    // Stream.generate(() -> atomic2.incrementAndGet()) // infinite, unordered, sequential
    // Stream.generate(() -> atomic2.incrementAndGet()).parallel() // infinite, unordered, parallel
}
Run Code Online (Sandbox Code Playgroud)

那我可以安全地对此流调用哪些操作以编写正确的算法?

看来,如果我可能想将元素写入文件中是一种副作用,那么我需要担心流是并行的:

// if stream is parallel, which order will be written to file?
coordinates().peek(i -> {writeToFile(i)}).count();
// how should I remember to always add sequential() in  such cases?
Run Code Online (Sandbox Code Playgroud)

而且,如果它是并行的,则基于什么Threadpool是并行的?

如果我想对流进行排序(或其他非短路操作),则在某种程度上需要谨慎对待它的无限性:

coordinates().sorted().limit(1000).collect(toList()); // will this terminate?
coordinates().allMatch(x -> x > 0); // will this terminate?
Run Code Online (Sandbox Code Playgroud)

我可以在排序之前强加一个限制,但是如果我期望一个未知大小的有限流,那应该是哪个幻数呢?

最后也许我想并行计算以节省时间,然后收集结果:

// will result list maintain the same order as sequential?
coordinates().map(i -> complexLookup(i)).parallel().collect(toList());
Run Code Online (Sandbox Code Playgroud)

但是,如果未对流进行排序(在该版本的库中),则由于并行处理,结果可能会混乱。但是除了不使用并行(这违反了性能目的)之外,我该如何防范呢?

集合明确表示是有限的还是无限的,是否有序,并且它们不带有处理模式或线程池。这些似乎是API的宝贵属性。

另外,有时可能需要关闭Streams,但最通常不需要关闭。如果我使用某个方法的流(来自某个方法参数的),通常应该调用close吗?

同样,流可能已经被消耗掉了,能够优雅地处理这种情况将是一个好习惯,因此检查流是否已经被消耗将是一个很好的选择。

我希望有一些代码片段可用于在处理之前验证有关流的假设,例如>

// if stream is parallel, which order will be written to file?
coordinates().peek(i -> {writeToFile(i)}).count();
// how should I remember to always add sequential() in  such cases?
Run Code Online (Sandbox Code Playgroud)

ori*_*rab 2

经过一番观察(一些实验和这里),没有办法确切地知道流是否是有限的。

不仅如此,有时甚至除了运行时之外都无法确定(例如在 java 11 中 -IntStream.generate(() -> 1).takeWhile(x -> externalCondition(x))),它也无法确定。

你能做的是:

  1. 您可以通过几种方式确定它是否是有限的(请注意,在这些方面接收 false 并不意味着它是无限的,只是可能如此):

    1. stream.spliterator().getExactSizeIfKnown()- 如果它的确切大小已知,则它是有限的,否则将返回 -1。

    2. stream.spliterator().hasCharacteristics(Spliterator.SIZED)- 如果是则SIZED返回 true。

  2. 您可以通过假设最坏的情况(取决于您的情况)来保护自己。

    1. stream.sequential()/stream.parallel()- 明确设置您的首选消费类型。
    2. 对于潜在的无限流,请假设每种情况下最坏的情况。

      1. 例如,假设您想收听一系列推文,直到找到Venkat的一条推文的推文 - 这可能是无限操作,但您希望等到找到这样的推文。因此,在这种情况下,只需继续stream.filter(tweet -> isByVenkat(tweet)).findAny()- 它将迭代,直到出现这样的推文(或永远)。
      2. 另一种情况(可能是更常见的情况)是想要对所有元素执行某些操作,或者仅尝试一定的时间(类似于超时)。为此,我建议stream.limit(x)在致电您的操作(collectallMatch类似操作)之前先致电x您愿意容忍的尝试次数。

毕竟,我只想提一下,我认为返回流通常不是一个好主意,除非有很大的好处,否则我会尽量避免它。