可以从仅暴露迭代器的"readNext"部分的对象创建流吗?

use*_*774 7 java java-8

我试图从csv文件读取,但由于它的大小,没有首先将它全部加载到内存中.

我找到的用于阅读csv的库是opencsv,它工作得非常好,但只暴露了两种方法:

readAll() 
Run Code Online (Sandbox Code Playgroud)

和

readNext() 
Run Code Online (Sandbox Code Playgroud)

readAll因为我不想在内存中同时使用所有内容,所以我想通过lazily从文件中读取来readNext.理想情况下,我想通过一个流来结束阅读.

我得到的最接近的是将readnext方法赋予Stream.generate构造,

Stream csvDataStream = Stream.generate(csvReader::readNext); 
Run Code Online (Sandbox Code Playgroud)

但是,一旦迭代器底层csvReader耗尽,这显然会导致抛出错误.我真的不想将我的整个程序包装在try/catch块中,因为我使用的是错误的语言.有没有办法从只暴露next方法的东西创建流?

Mar*_*nik 8

这是我项目中现成的实现.我有一个抽象的分裂器,它处理拆分成固定大小的批次,并允许有效地并行处理任何类型的基于I/O的流源:

import static java.util.Spliterators.spliterator;

import java.util.Comparator;
import java.util.Spliterator;
import java.util.function.Consumer;

public abstract class FixedBatchSpliteratorBase<T> implements Spliterator<T> {
  private final int batchSize;
  private final int characteristics;
  private long est;

  public FixedBatchSpliteratorBase(int characteristics, int batchSize, long est) {
    characteristics |= ORDERED;
    if ((characteristics & SIZED) != 0) characteristics |= SUBSIZED;
    this.characteristics = characteristics;
    this.batchSize = batchSize;
    this.est = est;
  }
  public FixedBatchSpliteratorBase(int characteristics, int batchSize) {
    this(characteristics, batchSize, Long.MAX_VALUE);
  }
  public FixedBatchSpliteratorBase(int characteristics) {
    this(characteristics, 64, Long.MAX_VALUE);
  }

  @Override public Spliterator<T> trySplit() {
    final HoldingConsumer<T> holder = new HoldingConsumer<>();
    if (!tryAdvance(holder)) return null;
    final Object[] a = new Object[batchSize];
    int j = 0;
    do a[j] = holder.value; while (++j < batchSize && tryAdvance(holder));
    if (est != Long.MAX_VALUE) est -= j;
    return spliterator(a, 0, j, characteristics());
  }
  @Override public Comparator<? super T> getComparator() {
    if (hasCharacteristics(SORTED)) return null;
    throw new IllegalStateException();
  }
  @Override public long estimateSize() { return est; }
  @Override public int characteristics() { return characteristics; }

  static final class HoldingConsumer<T> implements Consumer<T> {
    Object value;
    @Override public void accept(T value) { this.value = value; }
  }
}
Run Code Online (Sandbox Code Playgroud)

这是基于它的opencsv分裂器:

public class CsvSpliterator extends FixedBatchSpliteratorBase<String[]> {
  private final CSVReader cr;

  CsvSpliterator(CSVReader cr, int batchSize) {
    super(NONNULL, batchSize);
    if (cr == null) throw new NullPointerException("CSVReader is null");
    this.cr = cr;
  }
  public CsvSpliterator(CSVReader cr) { this(cr, 100); }

  @Override public void forEachRemaining(Consumer<? super String[]> action) {
    if (action == null) throw new NullPointerException();
    uncheckRun(() -> { for (String[] row; (row = cr.readNext()) != null;) action.accept(row); });
  }
  @Override public boolean tryAdvance(Consumer<? super String[]> action) {
    if (action == null) throw new NullPointerException();
    return uncheckCall(() -> {
      final String[] row = cr.readNext();
      if (row == null) return false;
      action.accept(row);
      return true;
    });
  }
}
Run Code Online (Sandbox Code Playgroud)

其中,uncheckRun和uncheckCall是

public static <T> T uncheckCall(Callable<T> callable) {
  try { return callable.call(); }
  catch (Exception e) { return sneakyThrow(e); }
}
public static void uncheckRun(RunnableExc r) {
  try { r.run(); } catch (Exception e) { sneakyThrow(e); }
}
public static <T> T sneakyThrow(Throwable e) {
  return Util.<RuntimeException, T>sneakyThrow0(e);
}
@SuppressWarnings("unchecked")
private static <E extends Throwable, T> T sneakyThrow0(Throwable t) throws E { throw (E)t; }
Run Code Online (Sandbox Code Playgroud)

用法:

import static java.util.stream.StreamSupport.stream;

....

final CSVReader cr = new CSVReader(new InputStreamReader(yourInputStream), separator, '"');
return stream(new CsvSpliterator(cr), true).onClose(() -> uncheckRun(cr::close));
Run Code Online (Sandbox Code Playgroud)


Bri*_*etz 7

实施一个Spliterator.您只需tryAdvance要用一个非常重要的实现来实现该方法; trySplit可以返回null,characteristics()可以返回ORDERED,并estimateSize可以返回Long.MAX_VALUE.然后打电话StreamSupport.stream(Spliterator)来制作一个流.