【问题标题】:Why stream created with Spliterators is not being processed in parallel?为什么使用 Spliterators 创建的流没有被并行处理?
【发布时间】:2020-10-20 12:31:50
【问题描述】:

这可能是非常基本的,但我不是 Java 人。这是我的处理代码,它只是打印和休眠:

    private static void myProcessings(int value)
    {
        System.out.println("Processing " + value);
    
        try
        {
            Thread.sleep(2000);
        }
        catch (InterruptedException e)
        {
            e.printStackTrace();
        }
    
        System.out.println("Finished processing " + value);
    }

现在,这个并行流似乎可以并行工作:

    IntStream iit = IntStream.rangeClosed(1,3);
    iit.parallel().forEach(Main::myProcessings);
    
    // output:

    // Processing 2
    // Processing 1
    // Processing 3
    // Finished processing 3
    // Finished processing 2
    // Finished processing 1

但是这个(由迭代器制成)没有:

    static class MyIter implements Iterator<Integer>
    {
        private int max;
        private int current;
    
        public MyIter(int maxVal)
        {
            max = maxVal;
            current = 1;
        }
    
        @Override
        public boolean hasNext()
        {
            return current <= max;
        }
    
        @Override
        public Integer next()
        {
            return current++;
        }
    }
    
    MyIter it = new MyIter(3);
    StreamSupport.stream(Spliterators.spliteratorUnknownSize(it, 0), true)
                 .forEach(Main::myProcessings);

    // output:

    // Processing 1
    // Finished processing 1
    // Processing 2
    // Finished processing 2
    // Processing 3
    // Finished processing 3

我在自定义迭代器版本中做错了什么? (我使用的是 Java 8)

【问题讨论】:

  • 你有parallel()在一个但没有在另一个?
  • @akuzminykh 后面的我也用过.parallel(),没用。因为StreamSupport.stream() 中的第二个参数已经使其并行。
  • @akuzminykh stream(Spliterators.spliteratorUnknownSize(it, 0), true) - 第二个参数 - true 是结果流是否必须并行的标志。
  • 相关且可能重复的stackoverflow.com/questions/46709455/…,它提到true 参数不会使流parallel
  • stackoverflow.com/a/48308511/4949750 - 检查这个答案 - 它准确地解释了这里的问题。 Spliterators 将仅使用足够大的集合(例如 10000 多个元素)来拆分工作。如果您减少睡眠时间并增加元素数量,您的代码就可以正常工作。

标签: java parallel-processing java-stream spliterator


【解决方案1】:

一种方法是估计流的大小:

Spliterators.spliterator(it, 3, 0);

数字(此处为 3)不必精确,但如果您给出 10000,则实际大小为 3 时只会使用一个线程。如果您给出 10,则会使用多个线程,甚至大小为 3。

估计值(在我的示例中为 3)用于确定批次的大小(在移动到下一个线程之前发送到一个线程的任务数)。如果您提供的估计数量很大并且只提交了几个任务,它们可能会全部分组并在第一个线程上运行,而不会向第二个线程发送任何内容。

【讨论】:

  • Iterator 不一定是线程安全的,这是问题的一部分。 Stream 必须在不支持并行处理的源上构建并行处理。
  • 问题是,我的对象有点笨重;比如说,存储在 db 中的图像,每个 80kB。我不能使用 10000,甚至 1000。我计划从 dB 获取 250 个对象;驱动程序为内存问题提供了一个迭代器。我有一个处理器函数来处理(CPU 绑定)这些对象。只是想使用 CPU 内核。
  • 看来,如果我有c cpu-cores,估计n 和实际数据数量m:如果n&gt;=m 它产生c 线程(或@ 987654329@,以较低者为准),这些线程完成所有处理。但是如果n&lt;m、c 线程将继续工作到c*n 任务。其余的m-c*n 任务将由主线程按顺序处理。所以对于n=0的问题中提到的情况;所有任务都按顺序处理。
  • @mshsayem 我的理解是该数字用于确定批次的大小(在移动到下一个线程之前发送到一个线程的任务数)。如果您提供的估计数量很大并且只提交了几个任务,那么它们可能都在同一个线程上运行。
  • 根据我的经验,在参考实现中,指定一个完全虚假的数字仍然比未知大小的流好得多。例如,偏离系数五通常是没有问题的。如果大小未知,您会遇到this comment 中描述的问题。即使在所有源元素都被缓冲之后,实现也不会理解它知道确切的大小。
猜你喜欢
  • 1970-01-01
  • 2015-06-03
  • 2015-10-23
  • 1970-01-01
  • 2014-12-18
  • 1970-01-01
  • 2020-10-29
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多