【问题标题】:Perform operation on stream in batches批量对流进行操作
【发布时间】:2018-02-14 20:12:05
【问题描述】:

我有 n 个条目的有限流。它的计数提前未知。数据大小大约为 10Gb,而且对于 RAM 来说太大了,所以我无法将其作为一个整体来阅读。 在每 100000 个条目之后分块处理该流的方法是什么?

Stream<?> blocks

所以我无法使用 subList 等方法。

我可以在代码中想象它是这样的:

    IntStream
            .range(0, Integer.MAX_VALUE)
            .filter(s -> s % 100000 == 0)
            .mapToObj(s -> blocks
                    .skip(s)
                    .collect(Collectors.toList())
                    .forEach(MyClass::doSomething)
            );

但是我得到了错误,因为 limit 是终端操作符并且它关闭了流。有什么解决方法吗?

【问题讨论】:

  • 我明白我为什么会收到这个错误,但这并不能帮助我想出解决方法。
  • 因为限制是终端操作员limit不是终端操作,你的代码中根本没有limit操作。
  • 另外,MyClass::doSomething 是什么?在.collect(Collectors.toList()) 之后链接.forEach(MyClass::doSomething) 表明无论如何都会为单个元素调用此方法,那么这里哪里需要额外的批处理呢? blocks.forEach(MyClass::doSomething) 足以一次处理一个元素,而无需一次将所有内容加载到内存中。

标签: java-stream


【解决方案1】:

如果

IntStream
        .range(0, Integer.MAX_VALUE)
        .filter(s -> s % 100000 == 0)
        .mapToObj(s -> blocks
                .skip(s)
                .collect(Collectors.toList())
                .forEach(MyClass::doSomething)
        );

编译没有错误,方法doSomething必须是一次接收单个元素blocks的方法,就像List.forEach(…)所做的那样,分别为每个元素调用消费者。 (忽略List.forEachvoid 的事实,因此对于外部流中的mapToObj(…) 是不够的)。在这种情况下,“批处理”根本没有任何好处,您可以使用

blocks.forEachOrdered(MyClass::doSomething);

因为这将一次加载一个元素,并允许每个元素在处理下一个元素时进行垃圾回收(除非doSomething 在某处存储引用)。

您在调用doSomething 之前尝试将100000 个元素收集到List 中并不能提高性能,因为流仍会一个接一个地加载每个元素,而您仍在处理一个接一个的元素。它只会阻止多达 99999 个元素的垃圾收集,直到第 100000 个元素被处理。这不是优势。

【讨论】:

    【解决方案2】:

    您似乎必须按照先前问题的答案中的建议使用 Spliterator have a look at the Java docs

    下面的简化示例将产生 10 个块块的输出,同时保持流打开。

    Stream<Block<Integer>> blocks = IntStream
      .range(0, 1000)
      .mapToObj(Block::new);
    
    Spliterator<Block<Integer>> split = blocks.spliterator();
    int chunkSize = 10;
    
    while (true) {
      List<Block<Integer>> chunk = new ArrayList<>(10);
      for (int i = 0; i < chunkSize && split.tryAdvance(chunk::add); i++) {
        chunk.get(i).doSomething();
      }
      System.out.println();
      if (chunk.isEmpty()) break;
    }
    
    class Block<T> {
      private T value;
    
      Block(T value) {
        this.value = value;
      }
    
      void doSomething() {
        System.out.print("v: " + value);
      }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2023-01-24
      • 2021-12-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多