【问题标题】:Java 8 parallel stream reduction - thread count vs speedupJava 8 并行流减少 - 线程数与加速
【发布时间】:2015-08-10 05:56:27
【问题描述】:

我正在针对 java 8 并行流的缩减操作进行不同线程数的实验:总计超过 1150 万个数字。我在 8 核英特尔至强处理器上运行它。当我将线程数从 16 更改为 1 时,我看不到总执行时间和总 CPU 时间有太大变化。我 确实得到了加速,而不是使用“流”的顺序版本的“平行流”。我想了解线程数量与顺序减少的加速之间的相关性。

有人可以帮我解释一下吗?我的代码在某处不正确吗?

源码为here,第74行进行并行归约。

代码的相关部分也粘贴在下面:

class RedOperator implements BinaryOperator<MonResult>{

    @Override
    public MonResult apply(MonResult t, MonResult u) {
        // TODO Auto-generated method stub
        if(t != null && u != null)
            t.count = t.count + u.count;
        return t;
    }

}

class FOFinisher implements Function<ConcurrentHashMap<ArrayList<String>, ArrayList<MonResult>>, ArrayList<MonResult>>{


    protected MonResult id;
    protected boolean isParallel;
    public FOFinisher(boolean isparallel){
        id = new MonResult();
        id.count = 0;
        this.isParallel = isparallel;
    }
    public long getJVMCpuTime() {
        long lastProcessCpuTime = 0;
        try {
            if (ManagementFactory.getOperatingSystemMXBean() instanceof OperatingSystemMXBean) {
                lastProcessCpuTime=((com.sun.management.OperatingSystemMXBean)ManagementFactory.getOperatingSystemMXBean()).getProcessCpuTime();
            }
        }
        catch (  ClassCastException e) {
            System.out.println(e.getMessage());
        }finally{
            return lastProcessCpuTime;
        }
    }
    @Override
    public ArrayList<MonResult> apply(ConcurrentHashMap<ArrayList<String>, ArrayList<MonResult>> t) {
        // TODO Auto-generated method stub

        long beg = System.nanoTime();
        long begCPU = getJVMCpuTime();
        ArrayList<MonResult> ret;
        RedOperator x = new RedOperator();
        if(isParallel){
            ret = new ArrayList<MonResult>(t.values().stream().map((alist)->{
                        return alist.parallelStream().reduce(id,x);
                    }).collect(Collectors.toList()));
        }
        else{
            ret = new ArrayList<MonResult>(t.values().stream().map((alist)->{
                        return alist.stream().reduce(id,x);
                    }).collect(Collectors.toList()));
        }

        long end = System.nanoTime();
        long endCPU = getJVMCpuTime();
        System.out.println("Exec Time:" + TimeUnit.MILLISECONDS.convert((end-beg),TimeUnit.NANOSECONDS));
        System.out.println("CPU Time : " + TimeUnit.MILLISECONDS.convert((endCPU-begCPU),TimeUnit.NANOSECONDS));
        return ret;
    }

}

【问题讨论】:

  • 请在问题中发布相关代码,不要链接。链接消失了,因此问题失去了所有上下文。
  • 如果您设法制作了一个独立的示例并在此处发布,我可以运行该示例来确认您的发现,那么我(和其他人)将能够为您提供帮助。
  • 您是否尝试过并行化外部流 (t.values().parallelStream())?此外,使用ret = t.values().&lt;...&gt;.collect(Collectors.toCollection(ArrayList::new)); 而不是new ArrayList&lt;&gt;(t.values()...collect(Collectors.toList()) 会更有效。

标签: java multithreading parallel-processing java-8 reduction


【解决方案1】:

与顺序版本相比,您看不到任何加速,因为大部分时间都花在了拆分和加入任务上。简而言之-您接受测试的问题太简单了。详情请查看Angelika Langer's Geecon presentation。

【讨论】:

    猜你喜欢
    • 2014-12-16
    • 2015-11-05
    • 1970-01-01
    • 2015-02-24
    • 1970-01-01
    • 1970-01-01
    • 2014-04-29
    • 2016-08-18
    • 2020-10-18
    相关资源
    最近更新 更多