Java Spliterator 不断拆分并行流

mar*_*ace 6 java java-stream

我发现 Java 并行流有一些令人惊讶的行为。我制作了自己的Spliterator,并且生成的并行流被分割,直到每个流中只有一个元素。这似乎太小了,我想知道我做错了什么。我希望我可以设置一些特征来纠正这个问题。

这是我的测试代码。在Float这里仅仅是一个虚拟的有效载荷,我真正的流类稍微复杂一些。

   public static void main( String[] args ) {
      TestingSpliterator splits = new TestingSpliterator( 10 );
      Stream<Float> test = StreamSupport.stream( splits, true );
      double total = test.mapToDouble( Float::doubleValue ).sum();
      System.out.println( "Total: " + total );
   }
Run Code Online (Sandbox Code Playgroud)

此代码将不断拆分此流,直到每个流Spliterator都只有一个元素。这似乎太多了,效率不高。

输出:

run:
Split on count: 10
Split on count: 5
Split on count: 3
Split on count: 5
Split on count: 2
Split on count: 2
Split on count: 3
Split on count: 2
Split on count: 2
Total: 5.164293184876442
BUILD SUCCESSFUL (total time: 0 seconds)
Run Code Online (Sandbox Code Playgroud)

这是Spliterator. 我主要关心的是我应该使用哪些特征,但也许其他地方有问题?

public class TestingSpliterator implements Spliterator<Float> {
   int count;
   int splits;

   public TestingSpliterator( int count ) {
      this.count = count;
   }

   @Override
   public boolean tryAdvance( Consumer<? super Float> cnsmr ) {
      if( count > 0 ) {
         cnsmr.accept( (float)Math.random() );
         count--;
         return true;
      } else
         return false;
   }

   @Override
   public Spliterator<Float> trySplit() {
      System.err.println( "Split on count: " + count );
      if( count > 1 ) {
         splits++;
         int half = count / 2;
         TestingSpliterator newSplit = new TestingSpliterator( count - half );
         count = half;
         return newSplit;
      } else
         return null;
   }

   @Override
   public long estimateSize() {
      return count;
   }

   @Override
   public int characteristics() {
      return IMMUTABLE | SIZED;
   }
}
Run Code Online (Sandbox Code Playgroud)

那么我怎样才能把流分成更大的块呢?我希望在 10,000 到 50,000 附近会更好。

我知道我可以null从该trySplit()方法返回,但这似乎是一种倒退的方法。系统似乎应该对内核数量、当前负载以及使用流的代码的复杂程度有所了解,并相应地调整自身。换句话说,我希望流块大小在外部配置,而不是由流本身在内部固定。

编辑:重新。Holger 在下面的回答中,当我增加原始流中的元素数量时,流拆分会稍微减少,因此StreamSupport最终会停止拆分。

初始流大小为 100 个元素时,StreamSupport当流大小达到 2 时停止拆分(我在屏幕上看到的最后一行是Split on count: 4)。

对于 1000 个元素的初始流大小,各个流块的最终大小约为 32 个元素。


编辑部分 deux:在查看了上面的输出后,我更改了我的代码以列出Spliterator创建的单个s。以下是变化:

   public static void main( String[] args ) {
      TestingSpliterator splits = new TestingSpliterator( 100 );
      Stream<Float> test = StreamSupport.stream( splits, true );
      double total = test.mapToDouble( Float::doubleValue ).sum();
      System.out.println( "Total Spliterators: " + testers.size() );
      for( TestingSpliterator t : testers ) {
         System.out.println( "Splits: " + t.splits );
      }
   } 
Run Code Online (Sandbox Code Playgroud)

TestingSpliterator's ctor:

   static Queue<TestingSpliterator> testers = new ConcurrentLinkedQueue<>();

   public TestingSpliterator( int count ) {
      this.count = count;
      testers.add( this ); // OUCH! 'this' escape
   }
Run Code Online (Sandbox Code Playgroud)

这段代码的结果是第一个Spliterator被拆分了 5 次。下一个Spliterator被拆分 4 次。下一组Spliteratorsget 拆分 3 次。等等。结果是制作了 36 个Spliterators,并且流被分成了尽可能多的部分。在典型的桌面系统上,这似乎是 API 认为最适合并行操作的方式。

我将在下面接受 Holger 的回答,这基本上是StreamSupport班级在做正确的事情,别担心,要开心。对我来说,部分问题是我正在对非常小的流大小进行早期测试,我对拆分的数量感到惊讶。不要自己犯同样的错误。

Hol*_*ger 3

你从错误的角度看待它。该实现并没有拆分 \xe2\x80\x9cuntil 每个 spliterator 有一个元素 \xe2\x80\x9d,而是拆分 \xe2\x80\x9cuntil 有 10 个 spliterator\xe2\x80\x9d。

\n

单个 spliterator 实例只能由一个线程处理。分裂器在遍历开始后不需要支持分裂。因此,任何事先未使用的分割机会都可能导致事后并行处理能力有限。

\n

重要的是要记住,Stream 实现收到了一个ToDoubleFunction未知工作负载的\xc2\xb9。它不知道它和你的情况一样简单Float::doubleValue。它可能是一个需要一分钟来评估的函数,然后,每个 CPU 核心都有一个分离器将是正确的。即使拥有多个 CPU 内核也是一种有效的策略,可以处理某些评估所需时间明显长于其他评估的可能性。

\n

初始分割器的典型数量为 \xe2\x80\x9cCPU 核心数\xe2\x80\x9d\xc2\xa0\xc3\x97\xc2\xa04,尽管稍后当更多地了解实际工作负载时,这里可能会有更多分割操作存在。当您的输入数据少于该数字时,当它被拆分直到每个拆分器留下一个元素时,它\xe2\x80\x99 并不奇怪。

\n

您可以尝试使用new TestingSpliterator( 10000 )1000100来查看,一旦实现假设有足够的块来保持所有 CPU 核心繁忙,分割数量不会发生显着变化。

\n

由于您的 spliterator 也不知道有关消费流的每个元素工作负载的任何信息,因此您不应该\xe2\x80\x99 担心这​​一点。如果您可以顺利地支持拆分为单个元素,那就这么做吧。

\n

\xc2\xb9 不过,\xe2\x80\x99 对于没有链接任何操作的情况没有特殊的优化。

\n

  • 正如[这个答案](/sf/answers/3136194911/)中所说,“*如果你请求并行,你就会得到并行,即使它实际上降低了性能。*”如图所示,你甚至可以处理当工作量足够大时,两个元素并行。请注意,分割的深度并不意味着它在不同的线程上运行;而是意味着它在不同的线程上运行。如果实际工作负载较低,则本地处理线程可能会在另一个线程拾取下一个块之前获取它。然后,您只需创建另一个轻量级对象的微小开销。“并行会有好处”没有神奇的门槛 (2认同)