【问题标题】:Nested Java 8 parallel forEach loop perform poor. Is this behavior expected?嵌套的 Java 8 并行 forEach 循环执行不佳。这种行为是预期的吗?
【发布时间】:2014-06-22 18:56:12
【问题描述】:

注意:我已经在另一篇 SO 帖子 - Using a semaphore inside a nested Java 8 parallel stream action may DEADLOCK. Is this a bug? - 中解决了这个问题,但这篇帖子的标题暗示这个问题与信号量的使用有关 - 这有点分散了讨论的注意力。我创建这个是为了强调嵌套循环可能存在性能问题——尽管这两个问题可能有一个共同的原因(也许是因为我花了很多时间来解决这个问题)。 (我不认为它是重复的,因为它强调了另一种症状 - 但如果你确实删除它)。

问题:如果嵌套两个Java 8 stream.parallel().forEach 循环并且所有任务都是独立的、无状态的等等——除了提交到公共FJ池——那么嵌套并行循环内的并行循环的性能比将顺序循环嵌套在并行循环内要差得多。更糟糕的是:如果包含内部循环的操作是同步的,你会得到一个死锁。

性能问题演示

没有“同步”,您仍然可以观察到性能问题。您可以在以下位置找到演示代码:http://svn.finmath.net/finmath%20experiments/trunk/src/net/finmath/experiments/concurrency/NestedParallelForEachTest.java (有关更详细的说明,请参阅那里的 JavaDoc)。

我们这里的设置如下:我们有一个嵌套的 stream.parallel().forEach()。

  • 内部循环是独立的(无状态、无干扰等 - 使用公共池除外)并且在最坏的情况下总共消耗 1 秒,即如果按顺序处理。
  • 外循环的一半任务在该循环之前消耗 10 秒。
  • 该循环后一半消耗 10 秒。
  • 因此每个线程总共消耗 11 秒(最坏情况)。 * 我们有一个布尔值,它允许将内部循环从并行()切换到顺序()。

现在:将 24 个外循环任务提交到并行度为 8 的池中,我们预计最多 24/8 * 11 = 33 秒(在 8 核或更好的机器上)。

结果是:

  • 使用内部顺序循环:33 秒。
  • 使用内部并行循环:>80 秒(我有 92 秒)。

问题:您能确认一下这种行为吗?这是人们对框架的期望吗? (我现在更加小心了,声称这是一个错误,但我个人认为这是由于 ForkJoinTask 的实现中的一个错误。备注:我已将此发布到并发兴趣(请参阅http://cs.oswego.edu/pipermail/concurrency-interest/2014-May/012652.html) ,但到目前为止我还没有从那里得到确认)。

僵局演示

下面的代码会死锁

    // Outer loop
    IntStream.range(0,numberOfTasksInOuterLoop).parallel().forEach(i -> {
        doWork();
        synchronized(this) {
            // Inner loop
            IntStream.range(0,numberOfTasksInInnerLoop).parallel().forEach(j -> {
                doWork();
            });
        }
    });

numberOfTasksInOuterLoop = 24numberOfTasksInInnerLoop = 240outerLoopOverheadFactor = 10000doWork 是一些无状态 CPU 刻录机。

您可以在 http://svn.finmath.net/finmath%20experiments/trunk/src/net/finmath/experiments/concurrency/NestedParallelForEachAndSynchronization.java 找到完整的演示代码 (有关更详细的说明,请参阅那里的 JavaDoc)。

这是预期的行为吗?请注意,有关 Java 并行流的文档没有提到任何嵌套或同步问题。此外,没有提到两者都使用共同的分叉连接池这一事实。

更新

另一个关于性能问题的测试可以在http://svn.finmath.net/finmath%20experiments/trunk/src/net/finmath/experiments/concurrency/NestedParallelForEachBenchmark.java找到 - 这个测试没有任何阻塞操作(没有 Thread.sleep 和不同步)。我在这里整理了一些评论:http://christian-fries.de/blog/files/2014-nested-java-8-parallel-foreach.html

更新 2

似乎这个问题和更严重的信号量死锁已经在 J​​ava8 u40 中得到修复。

【问题讨论】:

  • 也许值得链接到并发兴趣讨论:cs.oswego.edu/pipermail/concurrency-interest/2014-May/…
  • 那个并发兴趣讨论很奇怪。令我震惊的是,讨论中的所有成员都坚持这样的事情:“程序员应该知道 F/J 执行 xyz”或“程序员应该使用来自 F/ 的 for this”,但通过搜索entire documentation of java.util.stream,我发现根本没有提到 F/J 框架。我记得在某处看到过将 F/J 作为实现细节这样的提及,但那是完全不同的画面。
  • @Christian Fries 关于这个问题有任何消息吗?

标签: java concurrency parallel-processing java-8 java-stream


【解决方案1】:

我可以确认这仍然是 8u72 中的性能问题,尽管它不会再出现死锁。并行终端操作仍然使用ForkJoinPool 上下文之外的ForkJoinTask 实例完成,这意味着每个并行流仍然共享common pool

演示一个简单的病态案例:

import java.util.concurrent.ForkJoinPool;
import java.util.stream.IntStream;

public class ParallelPerf {

    private static final Object LOCK = new Object();

    private static void runInNewPool(Runnable task) {
        ForkJoinPool pool = new ForkJoinPool();
        try {
            pool.submit(task).join();
        } finally {
            pool.shutdown();
        }
    }

    private static <T> T runInNewPool(Callable<T> task) {
        ForkJoinPool pool = new ForkJoinPool();
        try {
            return pool.submit(task).join();
        } finally {
            pool.shutdown();
        }
    }

    private static void innerLoop() {
        IntStream.range(0, 32).parallel().forEach(i -> {
//          System.out.println(Thread.currentThread().getName());
            try {
                Thread.sleep(5);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        });
    }

    public static void main(String[] args) {
        System.out.println("==DEFAULT==");
        long startTime = System.nanoTime();
        IntStream.range(0, 32).parallel().forEach(i -> {
            synchronized (LOCK) {
                innerLoop();
            }
//          System.out.println(" outer: " + Thread.currentThread().getName());
        });
        System.out.println(System.nanoTime() - startTime);

        System.out.println("==NEW POOLS==");
        startTime = System.nanoTime();
        IntStream.range(0, 32).parallel().forEach(i -> {
            synchronized (LOCK) {
                runInNewPool(() -> innerLoop());
            }
//          System.out.println(" outer: " + Thread.currentThread().getName());
        });
        System.out.println(System.nanoTime() - startTime);
    }
}

第二次运行将innerLoop 传递给runInNewPool,而不是直接调用它。在我的机器上(i7-4790,8 个 CPU 线程),我得到了大约 4 倍的加速:

==DEFAULT==
4321223964
==NEW POOLS==
1015314802

取消注释其他打印语句会使问题变得明显:

[...]
ForkJoinPool.commonPool-worker-6
ForkJoinPool.commonPool-worker-6
ForkJoinPool.commonPool-worker-6
 outer: ForkJoinPool.commonPool-worker-6
ForkJoinPool.commonPool-worker-3
ForkJoinPool.commonPool-worker-3
[...]
ForkJoinPool.commonPool-worker-3
ForkJoinPool.commonPool-worker-3
 outer: ForkJoinPool.commonPool-worker-3
ForkJoinPool.commonPool-worker-4
ForkJoinPool.commonPool-worker-4
[...]

公共池工作线程堆积在同步块中,一次只能进入一个线程。由于内部并行操作使用同一个池,并且池中的所有其他线程都在等待锁,所以我们得到单线程执行。

以及使用单独的 ForkJoinPool 实例的结果:

[...]
ForkJoinPool-1-worker-0
ForkJoinPool-1-worker-6
ForkJoinPool-1-worker-5
 outer: ForkJoinPool.commonPool-worker-4
ForkJoinPool-2-worker-1
ForkJoinPool-2-worker-5
[...]
ForkJoinPool-2-worker-7
ForkJoinPool-2-worker-3
 outer: ForkJoinPool.commonPool-worker-1
ForkJoinPool-3-worker-2
ForkJoinPool-3-worker-5
[...]

我们仍然让内部循环一次在一个工作线程上运行,但内部并行操作每次都会获得一个新池,并且可以利用其所有工作线程。

这是一个人为的例子,但删除同步块仍然显示出类似的速度差异,因为内部和外部循环仍然在相同的工作线程上竞争。多线程应用程序在多个线程中使用并行流时需要小心,因为这可能会在它们重叠时导致随机减速。

这是所有终端操作的问题,不仅仅是forEach,因为它们都在公共池中运行任务。我使用上面的runInNewPool 方法作为解决方法,但希望这将在某个时候内置到标准库中。

【讨论】:

    【解决方案2】:

    稍微整理一下代码。我在 Java 8 update 45 上没有看到相同的结果。无疑会有开销,但与您所说的时间跨度相比,它非常小。

    可能会发生死锁,因为您正在使用外部循环消耗池中的所有可用线程,而没有剩余线程可以执行内部循环。

    以下程序打印

    isInnerStreamParallel: false, isCPUTimeBurned: false
    java.util.concurrent.ForkJoinPool.common.parallelism = 8
    Done in 33.1 seconds.
    isInnerStreamParallel: false, isCPUTimeBurned: true
    java.util.concurrent.ForkJoinPool.common.parallelism = 8
    Done in 33.0 seconds.
    isInnerStreamParallel: true, isCPUTimeBurned: false
    java.util.concurrent.ForkJoinPool.common.parallelism = 8
    Done in 32.5 seconds.
    isInnerStreamParallel: true, isCPUTimeBurned: true
    java.util.concurrent.ForkJoinPool.common.parallelism = 8
    Done in 32.6 seconds.
    

    代码

    import java.util.stream.IntStream;
    
    public class NestedParallelForEachTest {
        // Setup: Inner loop task 0.01 sec in worse case. Outer loop task: 10 sec + inner loop. This setup: (100 * 0.01 sec + 10 sec) * 24/8 = 33 sec.
        static final int numberOfTasksInOuterLoop = 24;  // In real applications this can be a large number (e.g. > 1000).
        static final int numberOfTasksInInnerLoop = 100;                // In real applications this can be a large number (e.g. > 1000).
        static final int concurrentExecutionsLimitForStreams    = 8;    // java.util.concurrent.ForkJoinPool.common.parallelism
    
        public static void main(String[] args) {
            testNestedLoops(false, false);
            testNestedLoops(false, true);
            testNestedLoops(true, false);
            testNestedLoops(true, true);
        }
    
        public static void testNestedLoops(boolean isInnerStreamParallel, boolean isCPUTimeBurned) {
            System.out.println("isInnerStreamParallel: " + isInnerStreamParallel + ", isCPUTimeBurned: " + isCPUTimeBurned);
            long start = System.nanoTime();
    
            System.setProperty("java.util.concurrent.ForkJoinPool.common.parallelism",Integer.toString(concurrentExecutionsLimitForStreams));
            System.out.println("java.util.concurrent.ForkJoinPool.common.parallelism = " + System.getProperty("java.util.concurrent.ForkJoinPool.common.parallelism"));
    
            // Outer loop
            IntStream.range(0, numberOfTasksInOuterLoop).parallel().forEach(i -> {
    //            System.out.println(i + "\t" + Thread.currentThread());
                if(i < 10) burnTime(10 * 1000, isCPUTimeBurned);
    
                IntStream range = IntStream.range(0, numberOfTasksInInnerLoop);
                if (isInnerStreamParallel) {
                    // Inner loop as parallel: worst case (sequential) it takes 10 * numberOfTasksInInnerLoop millis
                    range = range.parallel();
                } else {
                    // Inner loop as sequential
                }
                range.forEach(j -> burnTime(10, isCPUTimeBurned));
    
                if(i >= 10) burnTime(10 * 1000, isCPUTimeBurned);
            });
    
            long end = System.nanoTime();
    
            System.out.printf("Done in %.1f seconds.%n", (end - start) / 1e9);
        }
    
        static void burnTime(long millis, boolean isCPUTimeBurned) {
            if (isCPUTimeBurned) {
                long end = System.nanoTime() + millis * 1000000;
                while (System.nanoTime() < end)
                    ;
    
            } else {
                try {
                    Thread.sleep(millis);
                } catch (InterruptedException e) {
                    throw new AssertionError(e);
                }
            }
        }
    }
    

    【讨论】:

    • 似乎“bug”在 Java 8u40 之前已经悄然修复。 (我从未收到对我提交的错误报告的任何回复,因此我不知道在哪里可以找到更多信息)。
    【解决方案3】:

    问题是你配置的相当有限的并行度被外部流处理吃掉了:如果你说你想要八个线程并使用parallel()处理一个超过八个项目的流,它将创建八个工作线程并让他们处理物品。

    然后在您的消费者中,您正在使用parallel() 处理另一个流,但没有剩余的工作线程。由于工作线程被阻塞等待内部流处理结束,ForkJoinPool 必须创建新的工作线程,这违反了您配置的并行度。在我看来,它不会回收这些扩展线程,而是让它们在处理后立即死亡。因此,在您的内部处理中,会创建和处理新线程,这是一项昂贵的操作。

    您可能会将其视为一个缺陷,即启动线程不参与并行流处理的计算,而只是等待结果,但即使已修复,您仍然会遇到一个很难解决的一般问题(如果有的话)修复:

    每当工作线程与外部流项目的数量之间的比率较低时,实现会将它们全部用于外部流,因为它不知道流是外部流。因此,并行执行内部流请求的工作线程比可用的多。使用调用者线程来参与计算可以修复它,使其性能等于串行计算,但在这里获得并行执行的优势并不适用于固定数量的工作线程的概念。

    请注意,您在这里只是触及了这个问题的表面,因为您对项目的处理时间相当平衡。如果内部项和外部项的处理都出现分歧(与同一级别的项相比),问题将更加严重。


    更新:通过分析和查看代码,ForkJoinPool 确实 尝试使用等待线程进行“工作窃取”,但根据 @987654326 是否使用不同的代码@ 是工作线程或其他线程。结果,一个工作线程实际上在等待大约 80% 的时间并且几乎没有工作,而其他线程确实对计算做出了贡献……


    更新 2:为了完整起见,这里使用 cmets.xml 中描述的简单并行执行方法。由于它将每个项目排入队列,因此当单个项目的执行时间相当短时,预计会有很多开销。所以这不是一个复杂的解决方案,而是一个演示,它可以在没有太多魔法的情况下处理长时间运行的任务……

    import java.lang.reflect.UndeclaredThrowableException;
    import java.util.concurrent.*;
    import java.util.function.IntConsumer;
    import java.util.stream.Collectors;
    import java.util.stream.IntStream;
    
    public class NestedParallelForEachTest1 {
        static final boolean isInnerStreamParallel = true;
    
        // Setup: Inner loop task 0.01 sec in worse case. Outer loop task: 10 sec + inner loop. This setup: (100 * 0.01 sec + 10 sec) * 24/8 = 33 sec.
        static final int numberOfTasksInOuterLoop = 24;  // In real applications this can be a large number (e.g. > 1000).
        static final int numberOfTasksInInnerLoop = 100; // In real applications this can be a large number (e.g. > 1000).
        static final int concurrentExecutionsLimitForStreams = 8;
    
        public static void main(String[] args) throws InterruptedException, ExecutionException {
            System.out.println(System.getProperty("java.version")+" "+System.getProperty("java.home"));
            new NestedParallelForEachTest1().testNestedLoops();
            E.shutdown();
        }
    
        final static ThreadPoolExecutor E = new ThreadPoolExecutor(
            concurrentExecutionsLimitForStreams, concurrentExecutionsLimitForStreams,
            2, TimeUnit.MINUTES, new SynchronousQueue<>(), (r,e)->r.run() );
    
        public static void parallelForEach(IntStream s, IntConsumer c) {
            s.mapToObj(i->E.submit(()->c.accept(i))).collect(Collectors.toList())
             .forEach(NestedParallelForEachTest1::waitOrHelp);
        }
        static void waitOrHelp(Future f) {
            while(!f.isDone()) {
                Runnable r=E.getQueue().poll();
                if(r!=null) r.run();
            }
            try { f.get(); }
            catch(InterruptedException ex) { throw new RuntimeException(ex); }
            catch(ExecutionException eex) {
                Throwable t=eex.getCause();
                if(t instanceof RuntimeException) throw (RuntimeException)t;
                if(t instanceof Error) throw (Error)t;
                throw new UndeclaredThrowableException(t);
            }
        }
        public void testNestedLoops(NestedParallelForEachTest1 this) {
            long start = System.nanoTime();
            // Outer loop
            parallelForEach(IntStream.range(0,numberOfTasksInOuterLoop), i -> {
                if(i < 10) sleep(10 * 1000);
                if(isInnerStreamParallel) {
                    // Inner loop as parallel: worst case (sequential) it takes 10 * numberOfTasksInInnerLoop millis
                    parallelForEach(IntStream.range(0,numberOfTasksInInnerLoop), j -> sleep(10));
                }
                else {
                    // Inner loop as sequential
                    IntStream.range(0,numberOfTasksInInnerLoop).sequential().forEach(j -> sleep(10));
                }
                if(i >= 10) sleep(10 * 1000);
            });
            long end = System.nanoTime();
            System.out.println("Done in "+TimeUnit.NANOSECONDS.toSeconds(end-start)+" sec.");
        }
        static void sleep(int milli) {
            try {
                Thread.sleep(milli);
            } catch (InterruptedException ex) {
                throw new AssertionError(ex);
            }
        }
    }
    

    【讨论】:

    • 我相信创建 10 个额外线程(我正在打印线程)并不能解释时间差异。请注意,外部循环只有 24 个任务!而且:然后我不明白为什么将内部循环包装在一个线程中(并立即将其与调用线程连接)确实可以解决问题。 (参见其他 SO 帖子)。
    • 哦,好吧,通过查看代码,如果等待线程是工作线程,ForkJoinPool 确实 会尝试对计算做出贡献。但是这种帮助一定是低效的,以至于没有帮助的等待会更快。逆天。或者,好吧,我已经阅读了太多关于 ForkJoin 的负面消息,以至于现在我并不感到惊讶。
    • 我相信有一个错误。对于嵌套循环,hg.openjdk.java.net/jdk8u/jdk8u/jdk/file/6be37bafb11a/src/share/… 中第 401 行的测试出错,因为对于内部循环,很可能所有线程都是 ForkJoinWorkerThread 的实例。如果我将内部循环包装在一个新的 Thread 类中,问题就消失了。
    • @Christian Fries:该测试确实是一个关键点,但我已经看到这两种选择都会导致“工作窃取”的尝试,但使用完全不同的算法。结果不同——我更新了我的答案……
    • @Christian Fries:当我阅读您的声明时,我已经进行了这样的测试,即使用自定义 Thread 会有所作为。这就是为什么我说,当启动线程是非工作线程时,“工作窃取”似乎效果更好。顺便说一句,我使用一种天真的“将每个项目发布到ExecutorService”方法对替代并行处理实现进行了一些测试。尽管包装和排队每个项目开销,但使用 8 个线程时它在 32 秒内运行。因此,如果 F/J 与在几个小时内被黑客入侵的简单方法相比表现如此糟糕,那么它不应该是 Stream 后端。
    猜你喜欢
    • 1970-01-01
    • 2021-09-24
    • 1970-01-01
    • 2020-03-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多