【问题标题】:ForkJoinPool different latency with different style of same codeForkJoinPool 不同延迟与不同风格的相同代码
【发布时间】:2017-12-21 18:02:29
【问题描述】:

我试图将 paralleStream 与自定义 ForkJoin 池一起使用,该任务执行网络调用。当我使用以下样式时

pool.submit(() -> {
        ioDelays.parallelStream().forEach(n -> {
            induceRandomSleep(n);
        });
    }).get();

如果我循环并一一提交任务,所花费的时间几乎是 11 倍,如下所示:

for (final Integer num : ioDelays) {
        ForkJoinTask<Integer> task =  pool.submit(() -> {
            return induceRandomSleep(num);
        });
        tasks.add(task);
    }
    int count = 0;
    final List<Integer> returnVals = new ArrayList<>();
    tasks.forEach(task -> {
        try {
            returnVals.add(task.get());
        } catch (InterruptedException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        } catch (ExecutionException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }
    });

如果使用parallelStream,是否会涉及到ForkJoinPool.common?这是模拟上述两种样式的整个程序

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.ForkJoinTask;

public class FJTPExperiment {

    public static void main(String[] args) throws InterruptedException, ExecutionException {
        ForkJoinPool pool = new ForkJoinPool(200);

        List<Integer> ioDelays = new ArrayList<>();
        for (int i = 0; i <2000; i++) {
            ioDelays.add( (int)(300 *Math.random() + 200));
        }
        int originalCount = 0;
        for (Integer val : ioDelays) {
            originalCount += val;
        }
        System.out.println("Expected " + originalCount);
        System.out.println(Thread.currentThread().getName() + " ::::Number of threads in common pool :" + ForkJoinPool.getCommonPoolParallelism());


        long beginTimestamp = System.currentTimeMillis();
        pool.submit(() -> {
            ioDelays.parallelStream().forEach(n -> {
                induceRandomSleep(n);
            });
        }).get();
        long endTimestamp = System.currentTimeMillis();
        System.out.println("Took " + (endTimestamp - beginTimestamp) + " ms");


        List<ForkJoinTask<Integer>> tasks = new ArrayList<>();
        beginTimestamp = System.currentTimeMillis();
        for (final Integer num : ioDelays) {
            ForkJoinTask<Integer> task =  pool.submit(() -> {
                return induceRandomSleep(num);
            });
            tasks.add(task);
        }
        int count = 0;
        final List<Integer> returnVals = new ArrayList<>();
        tasks.forEach(task -> {
            try {
                returnVals.add(task.get());
            } catch (InterruptedException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            } catch (ExecutionException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            }
        });
        endTimestamp = System.currentTimeMillis();
        for (Integer val : returnVals) {
            count += val;
        }
        System.out.println("Count " + count);
        System.out.println("Took " + (endTimestamp - beginTimestamp) + " ms");
    }


    public static int induceRandomSleep(int sleepInterval) {
        System.out.println(Thread.currentThread().getName() + " ::::sleeping for " + sleepInterval + " ms");
        try {
            Thread.sleep(sleepInterval);
            return sleepInterval;
        } catch (InterruptedException e) {
            e.printStackTrace();
            return sleepInterval;
        }
    }
}

【问题讨论】:

  • 看起来这种将任务提交到您自己的 ForkJoinPool 的技巧并非在所有情况下都有效。请查看stackoverflow.com/questions/36947336/…
  • 谢谢你,但我能够验证是否使用了自定义线程池,根据我打印线程名称的日志,在我的情况下它使用的是自定义线程池跨度>
  • 主要任务提交到您的自定义池,但 parallelStream().forEach() 使用默认的 ForkJoinPool.common。您可以使用System.setProperty("java.util.concurrent.ForkJoinPool.common‌​.parallelism", "200"); 进行检查,在这种情况下,两者花费的时间将彼此接近
  • 我部分同意,根据日志,它坚持并行但仍未使用公共池。我尝试了一个类似的变体“parallelStream().limit(200).forEach”,它的工作方式与您提到的相同。担心的是,如果我设置系统级别属性,那么一切都会受到影响
  • 来自之前提供的链接:官方不支持使用自定义的 ForkJoinPool 进行流处理,并且在使用 forEach 时,使用默认的池并行度来确定流拆分器的行为。因此在 for 循环中提交任务是使用自定义 ForkJoinPool 的唯一可靠方法

标签: java multithreading threadpool


【解决方案1】:

我最终找到了问题的答案,有两部分:

1) 只有一个任务被提交到 ForkJoinPool 它是如何产生多个线程的?

查看JDK implementation,似乎在调用parallelStream 时,它会检查当前线程是否为ForkJoinWorkerThread,如果是,则将任务推送到客户ForkJoinPool 的队列,如果不是,则将其推送到ForkJoinPool.common。这也通过日志进行了验证。

2)如果它工作为什么它很慢?

它很慢,因为并行度不是来自自定义 ForkJoinPool 的并行度,而是来自 ForkJoinPool.common 的并行度,默认情况下限制为Number of CPU cores -1。 JDK 实现是hereLEAF_TARGET 是派生的here。如果这必须正常工作,那么应该有一个分支从自定义线程池的并行性中派生出LEAF_TARGET

【讨论】:

    猜你喜欢
    • 2020-07-13
    • 2016-07-30
    • 1970-01-01
    • 1970-01-01
    • 2018-01-29
    • 1970-01-01
    • 2014-03-10
    • 2011-10-18
    • 1970-01-01
    相关资源
    最近更新 更多