【问题标题】:Do sorted and distinct immediately process the stream?sorted 和 distinct 会立即处理流吗?
【发布时间】:2018-03-31 13:00:15
【问题描述】:

想象一下我有这样的东西:

Stream<Integer> stream = Stream.of(2,1,3,5,6,7,9,11,10)
            .distinct()
            .sorted();

distinct()sorted() 的 javadocs 都说它们是“有状态的中间操作”。这是否意味着流在内部将执行诸如创建哈希集、添加所有流值之类的操作,然后看到sorted() 会将这些值放入排序列表或排序集中?还是比这更聪明?

换句话说,.distinct().sorted() 会导致 java 遍历流两次,还是 java 会延迟直到执行终端操作(例如 .collect)?

【问题讨论】:

    标签: java java-stream


    【解决方案1】:

    您提出了一个加载的问题,暗示必须在两个备选方案之间进行选择。

    有状态的中间操作必须存储数据,在某些情况下,直到存储所有元素才能将元素传递到下游,但这并没有改变这项工作被推迟到终端操作完成的事实已经开始了。

    说它必须“遍历流两次”也是不正确的。有完全不同的遍历正在进行,例如在sorted() 的情况下,首先,遍历将被排序的内部缓冲区上的源填充,其次,遍历缓冲区。对于distinct(),在顺序处理中不会发生第二次遍历,内部的HashSet只是用来判断是否向下游传递一个元素。

    所以当你运行时

    Stream<Integer> stream = Stream.of(2,1,3,5,3)
        .peek(i -> System.out.println("source: "+i))
        .distinct()
        .peek(i -> System.out.println("distinct: "+i))
        .sorted()
        .peek(i -> System.out.println("sorted: "+i));
    System.out.println("commencing terminal operation");
    stream.forEachOrdered(i -> System.out.println("terminal: "+i));
    

    打印出来

    commencing terminal operation
    source: 2
    distinct: 2
    source: 1
    distinct: 1
    source: 3
    distinct: 3
    source: 5
    distinct: 5
    source: 3
    sorted: 1
    terminal: 1
    sorted: 2
    terminal: 2
    sorted: 3
    terminal: 3
    sorted: 5
    terminal: 5
    

    表明在终端操作开始之前没有发生任何事情,并且来自源的元素立即传递distinct() 操作(除非是重复的),而所有元素在传递到下游之前都在sorted() 操作中缓冲。

    可以进一步证明distinct()不需要遍历整个流:

    Stream.of(2,1,1,3,5,6,7,9,2,1,3,5,11,10)
        .peek(i -> System.out.println("source: "+i))
        .distinct()
        .peek(i -> System.out.println("distinct: "+i))
        .filter(i -> i>2)
        .findFirst().ifPresent(i -> System.out.println("found: "+i));
    

    打印

    source: 2
    distinct: 2
    source: 1
    distinct: 1
    source: 1
    source: 3
    distinct: 3
    found: 3
    

    正如Jose Da Silva’s answer 所解释和演示的,缓冲量可能会随着有序并行流的变化而变化,因为必须先调整部分结果,然后才能将它们传递给下游操作。

    由于这些操作在实际终端操作已知之前不会发生,因此可能比 OpenJDK 中当前发生的优化更多(但可能会在不同的实现或未来版本中发生)。例如。 sorted().toArray() 可以使用并返回相同的数组,或者sorted().findFirst() 可以变成min(),等等。

    【讨论】:

      【解决方案2】:

      根据 javadoc,distinctsorted 方法都是有状态的中间操作。

      StreamOps 对此操作有以下说明:

      有状态操作可能需要在产生结果之前处理整个输入。例如,在查看流的所有元素之前,无法通过对流进行排序产生任何结果。因此,在并行计算下,一些包含有状态中间操作的管道可能需要对数据进行多次传递,或者可能需要缓冲重要数据。

      但是流的收集只发生在终端操作中(例如toArraycollectforEach),这两个操作都在管道中处理,数据流过它。不过,需要注意的重要一点是这些操作的执行顺序,distinct() 方法的 javadoc 说:

      对于有序流,不同元素的选择是稳定的(对于重复元素,会保留在遇到顺序中最先出现的元素。)对于无序流,不做稳定性保证。


      对于顺序流,当这个流被排序时,唯一检查的元素是前一个元素,当没有排序时,内部使用HashSet,因此在sort之后执行distinct会产生更好的性能.

      (注意:正如 Eugene 所评论的,在这种连续的流中性能提升可能很小,特别是当代码很热时,但仍然避免创建额外的时间 HashSet

      这里你可以看到更多关于distinctsort的顺序:

      Java Streams: How to do an efficient "distinct and sort"?


      另一方面,对于并行流,doc 表示:

      在并行管道中保持 distinct() 的稳定性相对昂贵(要求操作充当完整的屏障,并具有大量缓冲开销),并且通常不需要稳定性。如果您的情况语义允许,使用无序流源(例如 generate(Supplier))或使用 BaseStream.unordered() 删除排序约束可能会显着提高并行管道中 distinct() 的执行效率。

      full barrier operation 表示:

      必须先执行所有上游操作,然后才能启动下游。 Stream API 中只有两个完整的屏障操作:.sorted()(每次)和 .distinct()(在有序并行情况下)。

      因此,当使用并行流时,相反的顺序通常会更好(只要当前流是无序的),即在sorted 之前使用distinct,因为 sorted 可以在不同时开始接收元素正在处理中。

      使用相反的顺序,首先排序(无序的并行流),然后使用 distinct,在两者中都设置了障碍,首先必须为sort 处理(流)所有元素,然后为distinct 处理所有元素。

      这是一个例子:

      Function<String, IntConsumer> process = name ->
              idx -> {
                  TimeUnit.SECONDS.sleep(ThreadLocalRandom
                          .current().nextInt(3)); // handle exception or use 
                                                  // LockSupport.parkNanos(..) sugested by Holger
                  System.out.println(name + idx);
              };
      

      下面的函数接收一个名字,并返回一个 int 消费者,它从 0-2 秒休眠,然后打印。

      IntStream.range(0, 8).parallel() // n > number of cores
              .unordered() // range generates ordered stream (not sorted)
              .peek(process.apply("B"))
              .distinct().peek(process.apply("D"))
              .sorted().peek(process.apply("S"))
              .toArray(); // terminal operation
      

      这将打印 B 和 D 的混合,然后是所有 S(distinct 中没有障碍)。

      如果你改变sorteddistinct的顺序:

              // ... rest
              .sorted().peek(process.apply("S"))
              .distinct().peek(process.apply("D"))
              // ... rest
      

      这将打印所有 B,然后是所有 S,然后是所有 D(distinct 中的障碍)。

      如果您想尝试更多,请在sorted 之后再次添加unordered

              // ... rest
              .sorted().unordered().peek(process.apply("S"))
              .distinct().peek(process.apply("D"))
              // ... rest
      

      这将打印所有 B,然后是 S 和 D 的混合(distinct 再次没有障碍)。


      编辑:

      将代码稍作更改,以便更好地解释和使用ThreadLocalRandom.current().nextInt(3),如建议的那样。

      【讨论】:

      • distinct() 不需要需要收集整个流。此外,您不应使用new Random().nextInt() % 3,因为结果可能会变为负数。有一个简单的专用方法来获取有界的 int:new Random().nextInt(3)。但更好的选择是ThreadLocalRandom.current().nextInt(3)。顺便说一句,Java 甚至还有一个不那么知名的方法,可以在不需要异常处理的情况下休眠:LockSupport.parkNanos(TimeUnit.SECONDS.toNanos(ThreadLocalRandom.current().nextInt(3)));
      • @JoseDaSilva 代码热时性能差异真的很小...stackoverflow.com/a/48624167/1059372
      • @Holger 关于 ramdom 你是对的,我改成ThreadLocalRandom,但那部分代码仍然不是最重要的。不知道LockSupport.parkNanos() 好提示。最后distinct 实际上确实需要在有序并行流上收集整个流(是一个屏障操作),就像在 javadoc 和代码中所说的那样。
      • @Eugene 是的,我在连续流的部分添加了一个注释,但即使性能提升微不足道,也存在。感谢评论。
      • @Holger * collect 不是合适的词,需要等待操作结束。
      猜你喜欢
      • 1970-01-01
      • 2010-11-07
      • 1970-01-01
      • 2018-05-29
      • 2019-05-28
      • 1970-01-01
      • 1970-01-01
      • 2015-09-14
      • 1970-01-01
      相关资源
      最近更新 更多