【问题标题】:Most efficient way to get the last element of a stream获取流的最后一个元素的最有效方法
【发布时间】:2014-12-18 13:17:25
【问题描述】:

Stream 没有last() 方法:

Stream<T> stream;
T last = stream.last(); // No such method

获取最后一个元素(或空流为 null)的最优雅和/或最有效的方法是什么?

【问题讨论】:

  • 如果您需要找到Stream 的最后一个元素,您可能需要重新考虑您的设计,并且如果您真的想使用StreamStreams 不一定是有序的或有限的。如果您的Stream 是无序的、无限的或两者兼有,则最后一个元素没有意义。在我看来,Stream 的意义在于在数据和处理方式之间提供一层抽象。因此,Stream 本身不需要了解其元素的相对顺序。找到Stream 中的最后一个元素是 O(n)。如果你有不同的数据结构,它可能是 O(1)。
  • @jeff 需求是真实的:情况大致是将商品添加到购物车,每次添加都返回错误信息(某些商品组合无效),但只有最后一次添加的错误信息(当所有项目都已添加,并且可以对购物车进行公平评估)是所需的信息。 (是的,我们使用的 API 已损坏,无法修复)。
  • @BrianGoetz:无限流也没有明确定义的 count(),但 Stream 仍然有 count() 方法。实际上,该论点适用于无限流上的任何非短路终端操作。
  • @BrianGoetz 我认为流应该有last() 方法。 4 月 1 日可能会有一项调查,应该如何定义无限流。我建议:“它永远不会返回,它至少 100% 使用一个处理器内核。在并行流上,它需要 100% 使用所有内核。”
  • 如果列表包含具有自然顺序或可以排序的对象,您可以使用max() 方法,如stream()...max(Comparator...)

标签: java java-8 java-stream


【解决方案1】:

做一个简单地返回当前值的归约:

Stream<T> stream;
T last = stream.reduce((a, b) -> b).orElse(null);

【讨论】:

  • 您认为这是优雅、高效还是两者兼而有之?
  • @Duncan 我认为两者兼而有之,但我还不是 Java 8 中的一把枪,前几天在工作中出现了这种需求 - 一个初级将流推入堆栈然后弹出它,并且我认为这看起来更好,但那里可能有更简单的东西。
  • 为了简洁和优雅,这个答案胜出。在一般情况下,它的效率相当高;它将相当好地并行化。对于一些知道其大小的流源,有一种更快的方法,但在大多数情况下,它不值得额外的代码来保存那几次迭代。
  • @BrianGoetz 这将如何很好地并行化?使用并行流将无法预测最后一个值
  • @BrianGoetz:仍然是O(n),即使除以 CPU 内核数。由于流不知道归约函数的作用,它仍然必须为每个元素评估它。
【解决方案2】:

这在很大程度上取决于Stream 的性质。请记住,“简单”并不一定意味着“高效”。如果您怀疑流非常大、执行繁重的操作或具有预先知道大小的源,则以下方法可能比简单的解决方案更有效:

static <T> T getLast(Stream<T> stream) {
    Spliterator<T> sp=stream.spliterator();
    if(sp.hasCharacteristics(Spliterator.SIZED|Spliterator.SUBSIZED)) {
        for(;;) {
            Spliterator<T> part=sp.trySplit();
            if(part==null) break;
            if(sp.getExactSizeIfKnown()==0) {
                sp=part;
                break;
            }
        }
    }
    T value=null;
    for(Iterator<T> it=recursive(sp); it.hasNext(); )
        value=it.next();
    return value;
}

private static <T> Iterator<T> recursive(Spliterator<T> sp) {
    Spliterator<T> prev=sp.trySplit();
    if(prev==null) return Spliterators.iterator(sp);
    Iterator<T> it=recursive(sp);
    if(it!=null && it.hasNext()) return it;
    return recursive(prev);
}

你可以用下面的例子来说明区别:

String s=getLast(
    IntStream.range(0, 10_000_000).mapToObj(i-> {
        System.out.println("potential heavy operation on "+i);
        return String.valueOf(i);
    }).parallel()
);
System.out.println(s);

它将打印:

potential heavy operation on 9999999
9999999

换句话说,它没有对前 9999999 个元素执行操作,而只对最后一个元素执行操作。

【讨论】:

  • hasCharacteristics() 块的意义何在?它增加了哪些recursive() 方法尚未涵盖的价值?后者已经导航到最后一个分割点。此外,recursive() 永远不会返回null,因此您可以删除it != null 检查。
  • 递归操作可以处理所有情况,但只是一种后备,因为它的递归深度与(未过滤的!)元素的数量相匹配。理想的情况是SUBSIZED 流,它可以保证非空分半,所以我们永远不需要回到左侧。请注意,在这种情况下,recursive 实际上不会递归,因为 trySplit 已经证明可以返回 null
  • 当然,代码可以用不同的方式编写,而且确实如此;我猜null-check 源于早期版本,但后来我发现对于非SUBSIZED 流,您必须处理可能的空拆分部分,即您必须迭代以确定它是否具有值,因此我将 Spliterators.iterator(…) 调用移动到 recursive 方法中,以便在右侧为空时能够备份到左侧。循环仍然是首选操作。
  • 有趣的解决方案。请注意,根据当前的 Stream API 实现,您的流必须是并行的或直接连接到源拆分器。否则,即使底层源拆分器拆分,它也会出于某种原因拒绝拆分。另一方面,您不能盲目地使用parallel(),因为这实际上可能会并行执行一些操作(如排序),意外地消耗更多的 CPU 内核。
  • @Tagir Valeev:对,示例代码使用了.parallel(),但实际上,它可以对sorted()distinct() 产生影响。我不认为,应该对任何其他中间操作产生影响……
【解决方案3】:

番石榴有Streams.findLast:

Stream<T> stream;
T last = Streams.findLast(stream);

【讨论】:

  • 而且它的性能比reduce((a, b) -&gt; b)好很多,因为它在内部使用Spliterator.trySplit
【解决方案4】:

这只是对Holger 答案的重构,因为代码虽然很棒,但有点难以阅读/理解,尤其是对于在 Java 之前不是 C 程序员的人。希望我重构的示例类对于那些不熟悉拆分器、它们做什么或它们如何工作的人来说更容易理解。

public class LastElementFinderExample {
    public static void main(String[] args){
        String s = getLast(
            LongStream.range(0, 10_000_000_000L).mapToObj(i-> {
                System.out.println("potential heavy operation on "+i);
                return String.valueOf(i);
            }).parallel()
        );
        System.out.println(s);
    }

    public static <T> T getLast(Stream<T> stream){
        Spliterator<T> sp = stream.spliterator();
        if(isSized(sp)) {
            sp = getLastSplit(sp);
        }
        return getIteratorLastValue(getLastIterator(sp));
    }

    private static boolean isSized(Spliterator<?> sp){
        return sp.hasCharacteristics(Spliterator.SIZED|Spliterator.SUBSIZED);
    }

    private static <T> Spliterator<T> getLastSplit(Spliterator<T> sp){
        return splitUntil(sp, s->s.getExactSizeIfKnown() == 0);
    }

    private static <T> Iterator<T> getLastIterator(Spliterator<T> sp) {
        return Spliterators.iterator(splitUntil(sp, null));
    }

    private static <T> T getIteratorLastValue(Iterator<T> it){
        T result = null;
        while (it.hasNext()){
            result = it.next();
        }
        return result;
    }

    private static <T> Spliterator<T> splitUntil(Spliterator<T> sp, Predicate<Spliterator<T>> condition){
        Spliterator<T> result = sp;
        for (Spliterator<T> part = sp.trySplit(); part != null; part = result.trySplit()){
            if (condition == null || condition.test(result)){
                result = part;
            }
        }
        return result;      
    }   
}

【讨论】:

    【解决方案5】:

    这是另一种解决方案(效率不高):

    List<String> list = Arrays.asList("abc","ab","cc");
    long count = list.stream().count();
    list.stream().skip(count-1).findFirst().ifPresent(System.out::println);
    

    【讨论】:

    • 有趣...你测试过这个吗?因为没有substream 方法,即使有这也行不通,因为count 是终端操作。那么这背后的故事是什么?
    • 奇怪,我不知道我有什么 jdk,但它确实有一个子流。我查看了官方的 javadoc(docs.oracle.com/javase/8/docs/api/java/util/stream/Stream.html),你是对的,它没有出现在这里。
    • 当然,您必须先检查count==0 是否作为Stream.skip 不喜欢-1 作为输入。除此之外,问题并没有说您可以两次获得Stream。也没有说两次获得Stream 就可以保证获得相同数量的元素。
    【解决方案6】:

    使用“skip”方法的并行无大小流很棘手,@Holger 的实现给出了错误的答案。 @Holger 的实现也有点慢,因为它使用了迭代器。

    @Holger 答案的优化:

    public static <T> Optional<T> last(Stream<? extends T> stream) {
        Objects.requireNonNull(stream, "stream");
    
        Spliterator<? extends T> spliterator = stream.spliterator();
        Spliterator<? extends T> lastSpliterator = spliterator;
    
        // Note that this method does not work very well with:
        // unsized parallel streams when used with skip methods.
        // on that cases it will answer Optional.empty.
    
        // Find the last spliterator with estimate size
        // Meaningfull only on unsized parallel streams
        if(spliterator.estimateSize() == Long.MAX_VALUE) {
            for (Spliterator<? extends T> prev = spliterator.trySplit(); prev != null; prev = spliterator.trySplit()) {
                lastSpliterator = prev;
            }
        }
    
        // Find the last spliterator on sized streams
        // Meaningfull only on parallel streams (note that unsized was transformed in sized)
        for (Spliterator<? extends T> prev = lastSpliterator.trySplit(); prev != null; prev = lastSpliterator.trySplit()) {
            if (lastSpliterator.estimateSize() == 0) {
                lastSpliterator = prev;
                break;
            }
        }
    
        // Find the last element of the last spliterator
        // Parallel streams only performs operation on one element
        AtomicReference<T> last = new AtomicReference<>();
        lastSpliterator.forEachRemaining(last::set);
    
        return Optional.ofNullable(last.get());
    }
    

    使用 junit 5 进行单元测试:

    @Test
    @DisplayName("last sequential sized")
    void last_sequential_sized() throws Exception {
        long expected = 10_000_000L;
        AtomicLong count = new AtomicLong();
        Stream<Long> stream = LongStream.rangeClosed(1, expected).boxed();
        stream = stream.skip(50_000).peek(num -> count.getAndIncrement());
    
        assertThat(Streams.last(stream)).hasValue(expected);
        assertThat(count).hasValue(9_950_000L);
    }
    
    @Test
    @DisplayName("last sequential unsized")
    void last_sequential_unsized() throws Exception {
        long expected = 10_000_000L;
        AtomicLong count = new AtomicLong();
        Stream<Long> stream = LongStream.rangeClosed(1, expected).boxed();
        stream = StreamSupport.stream(((Iterable<Long>) stream::iterator).spliterator(), stream.isParallel());
        stream = stream.skip(50_000).peek(num -> count.getAndIncrement());
    
        assertThat(Streams.last(stream)).hasValue(expected);
        assertThat(count).hasValue(9_950_000L);
    }
    
    @Test
    @DisplayName("last parallel sized")
    void last_parallel_sized() throws Exception {
        long expected = 10_000_000L;
        AtomicLong count = new AtomicLong();
        Stream<Long> stream = LongStream.rangeClosed(1, expected).boxed().parallel();
        stream = stream.skip(50_000).peek(num -> count.getAndIncrement());
    
        assertThat(Streams.last(stream)).hasValue(expected);
        assertThat(count).hasValue(1);
    }
    
    @Test
    @DisplayName("getLast parallel unsized")
    void last_parallel_unsized() throws Exception {
        long expected = 10_000_000L;
        AtomicLong count = new AtomicLong();
        Stream<Long> stream = LongStream.rangeClosed(1, expected).boxed().parallel();
        stream = StreamSupport.stream(((Iterable<Long>) stream::iterator).spliterator(), stream.isParallel());
        stream = stream.peek(num -> count.getAndIncrement());
    
        assertThat(Streams.last(stream)).hasValue(expected);
        assertThat(count).hasValue(1);
    }
    
    @Test
    @DisplayName("last parallel unsized with skip")
    void last_parallel_unsized_with_skip() throws Exception {
        long expected = 10_000_000L;
        AtomicLong count = new AtomicLong();
        Stream<Long> stream = LongStream.rangeClosed(1, expected).boxed().parallel();
        stream = StreamSupport.stream(((Iterable<Long>) stream::iterator).spliterator(), stream.isParallel());
        stream = stream.skip(50_000).peek(num -> count.getAndIncrement());
    
        // Unfortunately unsized parallel streams does not work very well with skip
        //assertThat(Streams.last(stream)).hasValue(expected);
        //assertThat(count).hasValue(1);
    
        // @Holger implementation gives wrong answer!!
        //assertThat(Streams.getLast(stream)).hasValue(9_950_000L); //!!!
        //assertThat(count).hasValue(1);
    
        // This is also not a very good answer better
        assertThat(Streams.last(stream)).isEmpty();
        assertThat(count).hasValue(0);
    }
    

    同时支持这两种情况的唯一解决方案是避免在未调整大小的并行流上检测到最后一个拆分器。结果是该解决方案将对所有元素执行操作,但它始终会给出正确的答案。

    请注意,在顺序流中,它无论如何都会对所有元素执行操作。

    public static <T> Optional<T> last(Stream<? extends T> stream) {
        Objects.requireNonNull(stream, "stream");
    
        Spliterator<? extends T> spliterator = stream.spliterator();
    
        // Find the last spliterator with estimate size (sized parallel streams)
        if(spliterator.hasCharacteristics(Spliterator.SIZED|Spliterator.SUBSIZED)) {
            // Find the last spliterator on sized streams (parallel streams)
            for (Spliterator<? extends T> prev = spliterator.trySplit(); prev != null; prev = spliterator.trySplit()) {
                if (spliterator.getExactSizeIfKnown() == 0) {
                    spliterator = prev;
                    break;
                }
            }
        }
    
        // Find the last element of the spliterator
        //AtomicReference<T> last = new AtomicReference<>();
        //spliterator.forEachRemaining(last::set);
    
        //return Optional.ofNullable(last.get());
    
        // A better one that supports native parallel streams
        return (Optional<T>) StreamSupport.stream(spliterator, stream.isParallel())
                .reduce((a, b) -> b);
    }
    

    关于该实现的单元测试,前三个测试完全相同(顺序和大小并行)。无大小并行的测试在这里:

    @Test
    @DisplayName("last parallel unsized")
    void last_parallel_unsized() throws Exception {
        long expected = 10_000_000L;
        AtomicLong count = new AtomicLong();
        Stream<Long> stream = LongStream.rangeClosed(1, expected).boxed().parallel();
        stream = StreamSupport.stream(((Iterable<Long>) stream::iterator).spliterator(), stream.isParallel());
        stream = stream.peek(num -> count.getAndIncrement());
    
        assertThat(Streams.last(stream)).hasValue(expected);
        assertThat(count).hasValue(10_000_000L);
    }
    
    @Test
    @DisplayName("last parallel unsized with skip")
    void last_parallel_unsized_with_skip() throws Exception {
        long expected = 10_000_000L;
        AtomicLong count = new AtomicLong();
        Stream<Long> stream = LongStream.rangeClosed(1, expected).boxed().parallel();
        stream = StreamSupport.stream(((Iterable<Long>) stream::iterator).spliterator(), stream.isParallel());
        stream = stream.skip(50_000).peek(num -> count.getAndIncrement());
    
        assertThat(Streams.last(stream)).hasValue(expected);
        assertThat(count).hasValue(9_950_000L);
    }
    

    【讨论】:

    • 请注意,单元测试使用 assertj 库以获得更好的流畅性。
    • 问题是你正在做StreamSupport.stream(((Iterable&lt;Long&gt;) stream::iterator).spliterator(), stream.isParallel()),经过一个完全没有特征的Iterable绕道,换句话说,创建了一个无序流。因此,结果与 parallel 或使用 skip 无关,只是因为“last”对于无序流没有意义,因此任何元素都是有效结果。跨度>
    【解决方案7】:

    我们在生产中需要 last 的 Stream - 我仍然不确定我们是否真的这样做了,但我团队中的多个团队成员说我们这样做是因为各种“原因”。我最终写了这样的东西:

     private static class Holder<T> implements Consumer<T> {
    
        T t = null;
        // needed to null elements that could be valid
        boolean set = false;
    
        @Override
        public void accept(T t) {
            this.t = t;
            set = true;
        }
    }
    
    /**
     * when a Stream is SUBSIZED, it means that all children (direct or not) are also SIZED and SUBSIZED;
     * meaning we know their size "always" no matter how many splits are there from the initial one.
     * <p>
     * when a Stream is SIZED, it means that we know it's current size, but nothing about it's "children",
     * a Set for example.
     */
    private static <T> Optional<Optional<T>> last(Stream<T> stream) {
    
        Spliterator<T> suffix = stream.spliterator();
        // nothing left to do here
        if (suffix.getExactSizeIfKnown() == 0) {
            return Optional.empty();
        }
    
        return Optional.of(Optional.ofNullable(compute(suffix, new Holder())));
    }
    
    
    private static <T> T compute(Spliterator<T> sp, Holder holder) {
    
        Spliterator<T> s;
        while (true) {
            Spliterator<T> prefix = sp.trySplit();
            // we can't split any further
            // BUT don't look at: prefix.getExactSizeIfKnown() == 0 because this
            // does not mean that suffix can't be split even more further down
            if (prefix == null) {
                s = sp;
                break;
            }
    
            // if prefix is known to have no elements, just drop it and continue with suffix
            if (prefix.getExactSizeIfKnown() == 0) {
                continue;
            }
    
            // if suffix has no elements, try to split prefix further
            if (sp.getExactSizeIfKnown() == 0) {
                sp = prefix;
            }
    
            // after a split, a stream that is not SUBSIZED can give birth to a spliterator that is
            if (sp.hasCharacteristics(Spliterator.SUBSIZED)) {
                return compute(sp, holder);
            } else {
                // if we don't know the known size of suffix or prefix, just try walk them individually
                // starting from suffix and see if we find our "last" there
                T suffixResult = compute(sp, holder);
                if (!holder.set) {
                    return compute(prefix, holder);
                }
                return suffixResult;
            }
    
    
        }
    
        s.forEachRemaining(holder::accept);
        // we control this, so that Holder::t is only T
        return (T) holder.t;
    
    }
    

    以及它的一些用法:

        Stream<Integer> st = Stream.concat(Stream.of(1, 2), Stream.empty());
        System.out.println(2 == last(st).get().get());
    
        st = Stream.concat(Stream.empty(), Stream.of(1, 2));
        System.out.println(2 == last(st).get().get());
    
        st = Stream.concat(Stream.iterate(0, i -> i + 1), Stream.of(1, 2, 3));
        System.out.println(3 == last(st).get().get());
    
        st = Stream.concat(Stream.iterate(0, i -> i + 1).limit(0), Stream.iterate(5, i -> i + 1).limit(3));
        System.out.println(7 == last(st).get().get());
    
        st = Stream.concat(Stream.iterate(5, i -> i + 1).limit(3), Stream.iterate(0, i -> i + 1).limit(0));
        System.out.println(7 == last(st).get().get());
    
        String s = last(
            IntStream.range(0, 10_000_000).mapToObj(i -> {
                System.out.println("potential heavy operation on " + i);
                return String.valueOf(i);
            }).parallel()
        ).get().get();
    
        System.out.println(s.equalsIgnoreCase("9999999"));
    
        st = Stream.empty();
        System.out.println(last(st).isEmpty());
    
        st = Stream.of(1, 2, 3, 4, null);
        System.out.println(last(st).get().isEmpty());
    
        st = Stream.of((Integer) null);
        System.out.println(last(st).isPresent());
    
        IntStream is = IntStream.range(0, 4).filter(i -> i != 3);
        System.out.println(last(is.boxed()));
    

    首先是Optional&lt;Optional&lt;T&gt;&gt; 的返回类型 - 它看起来很奇怪,我同意。如果第一个Optional 为空,则表示Stream 中没有元素;如果第二个 Optional 为空,则意味着最后一个元素实际上是 null,即:Stream.of(1, 2, 3, null)(不像 guavaStreams::findLast 在这种情况下会引发异常)。

    我承认我的灵感主要来自于 Holger 对我的一个类似问题的回答和 guava 的 Streams::findLast

    【讨论】:

      猜你喜欢
      • 2014-03-03
      • 2015-10-12
      • 1970-01-01
      • 2010-10-16
      • 2012-04-02
      • 1970-01-01
      • 2013-12-04
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多