与迭代器创建的Streams的并行性

Joh*_*han 4 java-8 java-stream

试验流我遇到了以下行为,我不太明白.我从迭代器创建了一个并行流,我注意到它似乎没有表现出并行性.在下面的例子中,我为控制台打印了一个计数器,用于两个并行流,一个是从迭代器创建的,另一个是从列表中创建的.从列表创建的流展示了我期望的行为,即以非顺序顺序打印计数器,但是从迭代器创建的流按顺序打印计数器.我是否错误地从迭代器创建并行流?

    private static int counter = 0;

public static void main(String[] args) {
    List<Integer> lstr = IntStream.rangeClosed(1, 100).boxed().collect(Collectors.toList());
    Iterator<Integer> iter = lstr.iterator();

    System.out.println("Iterator Stream: ");
    StreamSupport.stream(Spliterators.spliteratorUnknownSize(iter, Spliterator.IMMUTABLE | Spliterator.CONCURRENT), true).forEach(i -> {
        System.out.print(counter + " ");
        counter++;
    });

    counter = 0;
    System.out.println("\nList Stream: ");
    lstr.parallelStream().forEach(i -> {
        System.out.print(counter + " ");
        counter++;
    });

}
Run Code Online (Sandbox Code Playgroud)

Tag*_*eev 9

实施过程中存在缺陷Spliterators.spliteratorUnknownSize()。我在 Java 19 中修复了它,请参阅JDK-8280915。自 19-ea+19-1283 早期访问构建以来,该问题不再重现,并Spliterators.spliteratorUnknownSize已正确并行化。这是我机器上的输出:

Iterator Stream: 
0 0 0 1 0 0 6 7 8 0 0 0 0 12 0 15 0 17 18 10 20 21 22 23 11 0 5 27 28 29 30 0 31 32 3 35 36 37 38 39 40 41 42 4 0 2 45 45 43 34 33 26 25 24 54 19 56 16 0 59 59 0 62 14 13 13 65 64 63 61 60 57 71 73 73 75 76 55 54 79 52 51 50 49 48 47 46 85 84 80 74 70 70 68 67 66 94 92 88 98 
Run Code Online (Sandbox Code Playgroud)


Mis*_*sha 7

当前实现将尝试通过缓冲值并将它们分派给多个线程来并行化从迭代器生成的流,但只有在流足够长时它才会启动.将列表增加到10000个元素,您应该看到并行性.

使用大型列表,如果您收集到按线程分组的地图,则可能更容易看到线程分配的元素.更换你.forEach用.collect(Collectors.groupingBy(x -> Thread.currentThread().getName(), Collectors.counting()))


Hol*_*ger 7

没有保证并行处理以非连续顺序打印计数器.此外,由于您在没有同步的情况下更新变量,因此可能会错过其他线程所做的更新,因此结果可能完全不一致.

除此之外,Iterator必须按顺序轮询,因此要从并行处理中获得至少一些增益,必须缓冲元素,但是没有已知大小,对缓冲的元素数量没有很好的估计.默认策略使用超过一千个元素,并且不能很好地分割工作.

因此,如果您使用超过数千个元素,您可能会注意到更多并行活动.或者,您可以使用StreamSupport.stream(Spliterators.spliterator(iter, lstr.size(), 0), true)构造流来指定大小.然后,将调整内部使用的缓冲.

尽管如此,该List流将具有更高效的并行处理,因为它不仅知道其大小,还支持利用底层数据结构的随机访问性质来分割工作负载.

  • @Eugene`AbstractSpliterator`将首次拆分"1024"元素,每次将批量大小增加"1024",直到达到2²⁵.如果您的元素少于1000个,则所有元素都会出现在第一个数组中.问题是现在空的分裂器仍然报告"未知大小",但流实现将"未知大小"视为文字"Long.MAX_VALUE",因此它不再分割不到一千个元素数组,而是继续尝试拆分实际为空的<"unknow size"/"Long.MAX_VALUE"> spliterator.估计尺寸会导致它进行合理的分割. (4认同)
  • 难道他们都从1024开始,并通过加倍他们的大小增加?在电话里,现在不能看源:( (2认同)
  • `IteratorSpliterator` 与 `AbstractSpliterator` 的作用相同,但复制代码而不是继承它。 (2认同)