【问题标题】:Is there any way to stop a Stream.generate from its Lambda closure?有什么方法可以阻止 Stream.generate 从其 Lambda 闭包中停止?
【发布时间】:2014-03-25 09:47:48
【问题描述】:

我刚开始玩 Java 8 和 Lambda 表达式,我很好奇是否可以通过返回特定值来停止 Lambda 表达式内部的流生成 (如空)。 Stream.generate() 可以做到这一点吗?

private int counter;

private void generate()
{
    System.out.println(Stream.generate(() -> {
        if (counter < 10) {
            counter++;
            return RandomUtils.nextInt(100);
        } else {
            return null;
        }
    }).count());
}

很遗憾,此代码不会终止,因此仅返回 null 不会退出流。

【问题讨论】:

    标签: java lambda java-8


    【解决方案1】:

    Java 9 及更高版本包括this method

    Stream<T> takeWhile(Predicate<? super T> predicate); 
    

    按条件限制流。因此不再需要下面的解决方法。

    原始答案(适用于 9 之前的 Java 版本):

    对于 Stream.generate,这是根据定义不可能从 lambda 闭包中实现的。根据定义,它是无穷无尽的。使用limit(),您可以使您的流固定大小。但这对以下情况无济于事:

    if random>10 then stop
    

    有可能通过条件限制潜在的无限流。如果不知道大小,这很有用。您的朋友是 Spliterator,您的示例代码如下所示:

    System.out.println( StreamSupport.stream(Spliterators.spliteratorUnknownSize(new Iterator<Integer>() {
        int counter = 0;
    
        @Override
        public boolean hasNext() {
            return counter < 10;
        }
    
        @Override
        public Integer next() {
            counter++;
            return RandomUtils.nextInt(100);
        }
    }, Spliterator.IMMUTABLE), false).count());
    

    基本上,您可以从 Iterator 构建 Stream。我正在使用这个结构,例如对于来自 Stax XML 的 XMLEvents 流 - 解析。

    我知道这不是由 lambda 构造完成的,但它 IHMO 解决了按条件停止流项目生成这一缺少的功能。

    如果有更好的方法来实现这一点(我的意思是这种流构造而不是 XML 处理 ;))或者以这种方式使用流存在根本缺陷,我会非常感兴趣。

    【讨论】:

    • 有时无法知道是否有下一个元素,这使得两种解决方案都不可能。
    • 你能举个例子吗?
    • 想象一个流,它会扁平化在期货中返回的一组集合。您想在所有集合可用之前开始流式传输,因此您不知道未来集合中是否包含更多元素。
    【解决方案2】:

    这对于 Lamdas 是不可能的,您无法从表达式内部控制流程。 甚至 API 文档都说 Stream.generate 会生成一个无限流。

    但是,您可以通过使用 limit() 方法来限制 Stream 并实现所需的功能:

    System.out.println(Stream.generate(() -> RandomUtils.nextInt(100)).limit(10).count());
    

    【讨论】:

    • 这不是真的。
    【解决方案3】:
    // If you are not looking for parallelism, you can use following method:
    public static <T> Stream<T> breakStream(Stream<T> stream, Predicate<T> terminate) { 
      final Iterator<T> original = stream.iterator();
      Iterable<T> iter = () -> new Iterator<T>() { 
        T t;
        boolean hasValue = false;
    
        @Override
        public boolean hasNext() { 
          if (!original.hasNext()) { 
            return false;
          } 
          t = original.next();
          hasValue = true;
          if (terminate.test(t)) { 
            return false;
          } 
          return true;
        } 
    
        @Override
        public T next() { 
          if (hasValue) { 
            hasValue = false;
            return t;
          } 
          return t;
        } 
      };
    
      return StreamSupport.stream(iter.spliterator(), false);
    }
    

    【讨论】:

      【解决方案4】:

      使用StreamSupport.stream(Spliterator, boolean)
      请参阅Spliterator 上的 JavaDoc。
      这是示例拆分器:

      public class GeneratingSpliterator<T> implements Spliterator<T>
      {
          private Supplier<T> supplier;
          private Predicate<T> predicate;
      
          public GeneratingSpliterator(final Supplier<T> newSupplier, final Predicate<T> newPredicate)
          {
              supplier = newSupplier;
              predicate = newPredicate;
          }
      
          @Override
          public int characteristics()
          {
              return 0;
          }
      
          @Override
          public long estimateSize()
          {
              return Long.MAX_VALUE;
          }
      
          @Override
          public boolean tryAdvance(final Consumer<? super T> action)
          {
              T newObject = supplier.get();
              boolean ret = predicate.test(newObject);
              if(ret) action.accept(newObject);
              return ret;
          }
      
          @Override
          public Spliterator<T> trySplit()
          {
              return null;
          }
      }
      

      【讨论】:

        【解决方案5】:

        这是java 8的另一种解决方案(它需要一个Stream.Builder,可能不是最优的,但它很简单):

        @SuppressWarnings("ResultOfMethodCallIgnored")
        public static <T> Stream<T> streamBreakable(Stream<T> stream, Predicate<T> stopCondition) {
            Stream.Builder<T> builder = Stream.builder();
            stream.map(t -> {
                        boolean stop = stopCondition.test(t);
                        if (!stop) {
                            builder.add(t);
                        }
                        return stop;
                    })
                    .filter(result -> result)
                    .findFirst();
        
            return builder.build();
        }
        

        还有测试:

        @Test
        public void shouldStop() {
        
            AtomicInteger count = new AtomicInteger(0);
            Stream<Integer> stream = Stream.generate(() -> {
                if (count.getAndIncrement() < 10) {
                    return (int) (Math.random() * 100);
                } else {
                    return null;
                }
            });
        
            List<Integer> list = streamBreakable(stream, Objects::isNull)
                    .collect(Collectors.toList());
        
            System.out.println(list);
        }
        

        【讨论】:

          【解决方案6】:

          可能的,你只需要跳出框框思考。

          下面的想法是从 Python 中借来的,它是向我介绍生成器函数的语言...

          当你在 Supplier&lt;T&gt; 闭包中完成后,只需抛出一个 RuntimeException 的实例,然后在调用站点捕获并忽略它。

          一个示例摘录(请注意,我添加了Stream.limit(Long.MAX_VALUE) 的安全捕获以涵盖意外情况,尽管它不应该被触发):

          static <T> Stream<T> read(String path, FieldSetMapper<T> fieldSetMapper) throws IOException {
              ClassPathResource resource = new ClassPathResource(path);
              DefaultLineMapper<T> lineMapper = new DefaultLineMapper<>();
              lineMapper.setFieldSetMapper(fieldSetMapper);
              lineMapper.setLineTokenizer(getTokenizer(resource));
          
              return Stream.generate(new Supplier<T>() {
                  FlatFileItemReader<T> itemReader = new FlatFileItemReader<>();
                  int line = 1;
                  {
                      itemReader.setResource(resource);
                      itemReader.setLineMapper(lineMapper);
                      itemReader.setRecordSeparatorPolicy(new DefaultRecordSeparatorPolicy());
                      itemReader.setLinesToSkip(1);
                      itemReader.open(new ExecutionContext());
                  }
          
                  @Override
                  public T get() {
                      T item = null;
                      ++line;
                      try {
                          item = itemReader.read();
                          if (item == null) {
                              throw new StopIterationException();
                          }
                      } catch (StopIterationException ex) {
                          throw ex;
                      } catch (Exception ex) {
                          LOG.log(WARNING, ex,
                                  () -> format("%s reading line %d of %s", ex.getClass().getSimpleName(), line, resource));
                      }
                      return item;
                  }
              }).limit(Long.MAX_VALUE).filter(Objects::nonNull);
          }
          
          static class StopIterationException extends RuntimeException {}
          
          public void init() {
              if (repository.count() == 0) {
                  Level logLevel = INFO;
                  try {
                      read("providers.csv", fields -> new Provider(
                              fields.readString("code"),
                              fields.readString("name"),
                              LocalDate.parse(fields.readString("effectiveStart"), DateTimeFormatter.ISO_LOCAL_DATE),
                              LocalDate.parse(fields.readString("effectiveEnd"), DateTimeFormatter.ISO_LOCAL_DATE)
                      )).forEach(repository::save);
                  } catch (IOException e) {
                      logLevel = WARNING;
                      LOG.log(logLevel, "Initialization was interrupted");
                  } catch (StopIterationException ignored) {}
                  LOG.log(logLevel, "{} providers imported.", repository.count());
              }
          }
          

          【讨论】:

          • 这不是开箱即用的想法。那就是使用异常进行流量控制。由于多种原因,人们普遍认为这是一种不好的做法。
          • 你可以随心所欲地争论好的和坏的设计——这是唯一有效的选择。什么被认为是好或坏的一揽子规则带有警告,即用户应该了解分配状态的基本原理。在这种情况下,原因是潜在的代码混淆和异常处理设置的成本。这些原因在这里都无关紧要,并且它们被对工作解决方案的需求所抵消。在某些语言(如 Python)中,使用异常进行流控制是正常的 - 这是在核心语言类中迭代终止的方式。
          • 最终,不使用异常进行流控制源于异常是如何产生的——停止使用流控制(即魔术结果值)来表示异常状态。并允许异常处理与流控制正交。这些都是伟大的想法,但理想总是被实用主义压倒。在这种情况下,抛出异常是将消息传递给调用代码的唯一可用机制。
          【解决方案7】:

          我的解决方案是在完成后生成一个空值,然后应用过滤器

          Stream
           .generate( o -> newObject() )
           .filter( o -> o != null )
           .forEach(...)
          

          【讨论】:

          • 它不起作用,因为在这种情况下生成器不会停止。
          猜你喜欢
          • 2022-06-15
          • 1970-01-01
          • 2015-01-28
          • 2011-04-08
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2011-04-25
          • 1970-01-01
          相关资源
          最近更新 更多