【发布时间】: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().<...>.collect(Collectors.toCollection(ArrayList::new));而不是new ArrayList<>(t.values()...collect(Collectors.toList())会更有效。
标签: java multithreading parallel-processing java-8 reduction