【问题标题】:Implementing custom intermediate operations on Java 8 Streams在 Java 8 Streams 上实现自定义中间操作
【发布时间】:2019-06-16 04:07:42
【问题描述】:

我正在尝试研究如何在 Java 8 Stream 上实现自定义中间操作。看来我被锁定了:(

具体来说,我想获取一个流并将每个条目返回到并包括第一个具有特定值的条目。我想在那之后停止产生任何东西——让它短路。

它正在对输入数据进行一系列验证检查。我想在第一个错误上停下来,如果有的话,但我想在途中整理警告。而且因为这些验证检查可能很昂贵 - 例如涉及数据库查找 - 我只想运行所需的最小集合。

所以代码应该是这样的:

Optional<ValidationResult> result = validators.stream()
    .map(validator -> validator.validate(data))
    .takeUntil(result -> result.isError()) // This is the bit I can't do
    .reduce(new ValidationResult(), ::mergeResults);

似乎我应该能够用 ReferencePipeline.StatefulOp 做点什么,除了它都是包范围,所以我不能扩展它。所以我想知道实现这一目标的正确方法是什么?或者是否有可能?

还要注意 - 这需要在 Java 8 中,而不是 9+,因为由于各种不相关的原因我们还没有。

干杯

【问题讨论】:

  • 在 Java-9 中寻找takeWhile
  • 也许this 的答案会帮助您在Java 8 中创建自己的takeWhile()
  • ValidationResult 中的值是什么?是否可以忽略它的字段而只关心isError 以及剩下哪些验证器?如果是这样,请检查我的答案...
  • stackoverflow.com/questions/32290278/… 这似乎是您想要实现的目标。
  • @Naman takeWhile 无法工作,因为问题的包括部分

标签: java java-stream


【解决方案1】:

一般情况下,自定义操作需要处理Spliterator接口。它扩展了Iterator 的概念,通过添加特征和大小信息以及将部分元素拆分为另一个拆分器的能力(因此得名)。它还简化了迭代逻辑,只需要一种方法。

public static <T> Stream<T> takeWhile(Stream<T> s, Predicate<? super T> condition) {
    boolean parallel = s.isParallel();
    Spliterator<T> spliterator = s.spliterator();
    return StreamSupport.stream(new Spliterators.AbstractSpliterator<T>(
        spliterator.estimateSize(),
        spliterator.characteristics()&~(Spliterator.SIZED|Spliterator.SUBSIZED)) {
            boolean active = true;
            Consumer<? super T> current;
            Consumer<T> adapter = t -> {
                if((active = condition.test(t))) current.accept(t);
            };

            @Override
            public boolean tryAdvance(Consumer<? super T> action) {
                if(!active) return false;
                current = action;
                try {
                    return spliterator.tryAdvance(adapter) && active;
                }
                finally {
                    current = null;
                }
            }
        }, parallel).onClose(s::close);
}

为了保持流的属性,我们首先查询并行状态,为新流重新建立它。此外,我们还注册了一个关闭操作来关闭原始流。

主要工作是实现一个Spliterator来装饰之前流状态的拆分器。

除了SIZEDSUBSIZED 之外的特征被保留,因为我们的操作会导致不可预测的大小。原来的尺寸仍然是通过的,它现在将被用作估计值。

此解决方案在操作期间存储传递给tryAdvanceConsumer,以便能够使用相同的适配器消费者,避免为每次迭代创建一个新的消费者。这是可行的,因为它保证不会同时调用tryAdvance

并行是通过拆分完成的,它继承自AbstractSpliterator。这种继承的实现会缓冲一些元素,这是合理的,因为像takeWhile这样的操作实现更好的策略确实很复杂。

所以你可以像这样使用它

    takeWhile(Stream.of("foo", "bar", "baz", "hello", "world"), s -> s.length() == 3)
        .forEach(System.out::println);

将打印出来

foo
bar
baz

takeWhile(Stream.of("foo", "bar", "baz", "hello", "world")
    .peek(s -> System.out.println("before takeWhile: "+s)), s -> s.length() == 3)
    .peek(s -> System.out.println("after takeWhile: "+s))
    .forEach(System.out::println);

将打印出来

before takeWhile: foo
after takeWhile: foo
foo
before takeWhile: bar
after takeWhile: bar
bar
before takeWhile: baz
after takeWhile: baz
baz
before takeWhile: hello

这表明它不会处理超出必要的范围。在takeWhile阶段之前,我们必须遇到第一个不匹配的元素,之后,我们只会遇到直到那个元素。

【讨论】:

    【解决方案2】:

    我承认代码方面,Holger 的回答更性感,但可能更容易阅读:

    public static <T> Stream<T> takeUntilIncluding(Stream<T> s, Predicate<? super T> condition) {
    
        class Box implements Consumer<T> {
    
            boolean stop = false;
    
            T t;
    
            @Override
            public void accept(T t) {
                this.t = t;
            }
        }
    
        Box box = new Box();
    
        Spliterator<T> original = s.spliterator();
    
        return StreamSupport.stream(new AbstractSpliterator<>(
            original.estimateSize(),
            original.characteristics() & ~(Spliterator.SIZED | Spliterator.SUBSIZED)) {
    
            @Override
            public boolean tryAdvance(Consumer<? super T> action) {
    
                if (!box.stop && original.tryAdvance(box) && condition.test(box.t)) {
                    action.accept(box.t);
                    return true;
                }
    
                box.stop = true;
    
                return false;
            }
        }, s.isParallel());
    
    }
    

    【讨论】:

      【解决方案3】:

      你可以用一个技巧来做到这一点:

      List<ValidationResult> res = new ArrayList<>(); // Can modify it with your `mergeResults` instead of list
      
      Optional<ValidationResult> result = validators.stream()
          .map(validator -> validator.validate(data))
          .map(v -> {
             res.add(v);
             return v;
          })
          .filter(result -> result.isError())
          .findFirst();
      

      List&lt;ValidationResult&gt; res 将包含您感兴趣的值。

      【讨论】:

        【解决方案4】:

        可以使用以下结构;

        AtomicBoolean gateKeeper = new AtomicBoolean(true);    
        Optional<Foo> result = validators.stream()
            .filter(validator -> gateKeeper.get() 
                        && gateKeeper.compareAndSet(true, !validator.validate(data).isError()) 
                        && gateKeeper.get())
            .reduce(...) //have the first n non-error validators here
        

        带有gateKeeper 的过滤器充当短路逻辑并继续运行,直到遇到第一个isError() == true 案例,拒绝它,然后关闭其他validate() 呼叫的大门。它看起来有点疯狂,但它比其他 custom 实现简单得多,如果它符合您的要求,它可能会完美运行。

        不是 100% 确定这是否有用,因为我忽略了 validator.validate(data) 的结果,除了 isError() 结果,以及它属于列表中的 validator 的事实。

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 2019-07-09
          • 1970-01-01
          • 1970-01-01
          • 2018-05-07
          • 1970-01-01
          • 2023-03-04
          • 1970-01-01
          相关资源
          最近更新 更多