标签: parallel-processing

并行化 std::nth_element 和 std::partition

我正在将使用std::nth_element和的 C++ 代码移植std::partition到 OpenCL。

nth_element是一种选择算法,它将数组中第 n 个最小的数字放在第 n 个位置,并排列剩余元素,使所有小于该数字的元素在数组中位于该数字之前,所有大于该数字的元素位于该数字之后。实际上,nth_element将数组排序为 3 个桶:数字本身、所有小于该数字的数字以及所有大于该数字的数字。

规范地,nth_element是使用递归分区来实现的:选择一个元素,根据元素是否小于该元素来对元素进行分区。然后,选择包含数组第 n 个元素的存储桶并在该存储桶上递归。与完整快速排序之间的主要区别nth_element在于,快速排序在两个存储桶上递归,而不仅仅是包含第 n 个元素的存储桶。


partition是一个较弱的版本,nth_element它仅将数组分为 2 个桶:条件为 true 的桶和条件为 false 的桶。我链接到的网站给出了实现:

while (first!=last) {
    while (pred(*first)) {
        ++first;
        if (first==last) return first;
    }
    do {
        --last;
        if (first==last) return first;
    } while (!pred(*last));
    swap (*first,*last);
    ++first;
}
return first;
Run Code Online (Sandbox Code Playgroud)

其中 pred 是一个函数,用于评估某个元素是否应该位于第一个存储桶中。基本上,这个函数迭代地找到数组中位于错误位置的最外层元素对,并交换它们,当这对元素是相同元素时停止。


以下是我对并行化nth_element和的初步想法partition

分区可以使用原子比较和交换来实现,但我不确定如何覆盖所有可能交换的值对。没有明显的方法可以在多个线程之间划分工作,因为分区需要比较可能彼此相邻或位于数组两端的元素。我也没有找到一种方法来避免线程 B 与已被线程 A 交换的元素进行比较,这是低效的。 …

c++ sorting algorithm parallel-processing opencl

5
推荐指数
1
解决办法
2727
查看次数

在函数内部调用 clusterApply 时,性能会下降

我遇到了一个奇怪的问题clusterApply,我已经能够尽可能地隔离该问题,如下所示。首先,我从全局环境运行以下代码:

require(parallel)
cl<-makeCluster(rep("localhost",20),"SOCK")
xl<-list()
for(i in 1:20)
  xl[[i]]<-crossprod(matrix(rnorm(1e6),1000,1000))
x<-xl
clusterExport(cl,"x",environment())
f0<-function(z) eigen(x[[z]])
system.time(clusterApply(cl,1:20,f0))
##    user  system elapsed 
##   0.332   0.264   3.334 
Run Code Online (Sandbox Code Playgroud)

现在,为了确保没有发生任何奇怪的情况,请重新启动 R,然后运行以下类似的代码,该代码clusterApply从函数内部调用:

require(parallel)
cl<-makeCluster(rep("localhost",20),"SOCK")
xl<-list()
for(i in 1:20)
  xl[[i]]<-crossprod(matrix(rnorm(1e6),1000,1000))
f<-function(clust,x){
  force(x)
  clusterExport(clust,"x",environment())
  f0<-function(z) eigen(x[[z]])
  print(system.time(clusterApply(clust,1:20,f0)))
}
f(cl,xl)
##   user  system elapsed 
##  5.212   1.888  13.627 
Run Code Online (Sandbox Code Playgroud)

我做了一些搜索,找到了相关问题的答案,它指出未在全局环境中定义的函数中使用的局部变量将导出到集群。所以我想,也许问题在于x导出了两次,这才是花费很长时间的原因,而不是实际的函数调用。为了测试这一点,我将函数定义更改为:

f0<-function(z) eigen(get("x")[[z]])
Run Code Online (Sandbox Code Playgroud)

我的表现仍然很慢。有谁知道这里会发生什么?

顺便说一句,如果我只是打电话

clusterApply(clust,x,eigen)
Run Code Online (Sandbox Code Playgroud)

在函数内部,那么它就可以正常工作,就像在全局环境中一样快。当然,如果这是我想要解决的问题,我会简单地这样做,但事实并非如此,这只是一个玩具问题,用于隔离我与其他更复杂的代码遇到的问题。

parallel-processing r

5
推荐指数
1
解决办法
1307
查看次数

并行计算时如何写出日志?如何调试并行计算?

我发现如果并行计算期间有多个打印函数,则只有最后一个会显示在控制台上。所以我设置了outfile选项,希望我能得到每次打印的结果。这是 R 代码:

cl <- makeCluster(3, type = "SOCK",outfile="log.txt") 

abc <<- 123

clusterExport(cl,"abc")

clusterApplyLB(cl, 1:6,  
         function(y){
                     print(paste("before:",abc));
                     abc<<-y;
                     print(paste("after:",abc));
         }
)
stopCluster(cl)
Run Code Online (Sandbox Code Playgroud)

但我只得到三个记录:

starting worker for localhost:11888 
Type: EXEC 
Type: EXEC 
[1] "index: 3"
[1] "before: 123"
[1] "after: 2"
Type: EXEC 
[1] "index: 6"
[1] "before: 2"
[1] "after: 6"
Type: DONE 
Run Code Online (Sandbox Code Playgroud)

debugging parallel-processing r

5
推荐指数
1
解决办法
4797
查看次数

多个进程写入同一个CSV文件,如何避免冲突?

在我们的系统中,9 个进程同时写入相同的 CSV 输出。而且输出速度快。每天大约有 1000 万个新行。为了编写CSV文件,我们使用Python2.7的csv模块。

最近我注意到 CSV 文件中有一些混合行(参见下面的示例)。

例如

"name", "sex", "country", "email"
...# skip some lines
"qi", "Male", "China", "redice
...# skip some lines
"Jamp", "Male", "China", "jamp@site-digger.com"
...# skip some lines
@163.com"
Run Code Online (Sandbox Code Playgroud)

正确的输出应该是:

"name", "sex", "country", "email"
...# skip some lines
"qi", "Male", "China", "redice@163.com"
...# skip some lines
"Jamp", "Male", "China", "jamp@site-digger.com"
...
Run Code Online (Sandbox Code Playgroud)

如何避免这样的冲突呢?

python csv parallel-processing

5
推荐指数
1
解决办法
5222
查看次数

R Parallel,每次使用并行应用时创建一个新集群是否更好?

我正在 20 个核心上使用parLapply函数。我想其他功能也是一样的parSapply......

首先,将簇作为参数传递给函数以便该函数可以在不同的子函数之间调度簇的使用是一种不好的做法吗?

其次,我将此集群参数传递给一个函数,因此我认为每次使用时它都是同一个集群parLapply,每次调用都使用新集群会更好吗parLapply

谢谢

参考值

parallel-processing r

5
推荐指数
1
解决办法
340
查看次数

如何在类中并行化 python 中的 for ?

我有一个 python 函数funz,每次都会返回长度为 p 的不同数组。我需要多次运行该函数,然后计算每个值的平均值。

我可以使用 for 循环来完成此操作,但需要很多次。

我正在尝试使用库多处理,但遇到错误。

import sklearn as sk
import numpy as np
from sklearn.base import BaseEstimator, TransformerMixin
from sklearn import preprocessing,linear_model, cross_validation
from scipy import stats
from multiprocessing import Pool


class stabilize(BaseEstimator,TransformerMixin):

    def __init__(self,sim=3,n_folds=3):
        self.sim=sim
        self.n_folds=n_folds

    def fit(self,X,y):
        self.n,self.p=X.shape
        self.X=X
        self.y=y        
        self.beta=np.zeros(shape=(self.sim,self.p))
        self.alpha_min=[]        
        self.mapper=p.map(self.multiple_cv,[1]*self.sim)    

    def multiple_cv(self,o):
        kf=sk.cross_validation.KFold(self.n,n_folds=self.n_folds,shuffle=True)
        cv=sk.linear_model.LassoCV(cv=kf).fit(self.X,self.y)
        beta=cv.coef_
        alpha_min=cv.alpha_
        return alpha_min
Run Code Online (Sandbox Code Playgroud)

我使用虚拟变量 o 来告诉我要使用多少个并行进程。这不是很优雅,也许是错误的一部分。变量 X 和 y 已经是类的一部分,因此我没有参数传递给函数 multiple_cv。

当我运行该程序时,我收到此错误

Exception in thread Thread-3:
Traceback (most recent call last):
  File "/usr/lib/python2.7/threading.py", line …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing pool multiprocessing

5
推荐指数
1
解决办法
2589
查看次数

使用 Parallel.ForEach() 时防止远程系统过载

我们构建了这个应用程序,需要在远程计算机(实际上是 MatLab 服务器)上完成一些计算。我们使用 Web 服务连接到 MatLab 服务器并执行计算。

为了加快速度,我们使用了Parallel.ForEach()同时进行多个服务调用的方法。如果我们非常保守地将ParallelOptions.MaxDegreeOfParallelism(DOP) 设置为 4 或其他值,那么一切都会运行良好。然而,如果我们让框架决定 DOP,它将产生如此多的线程,从而迫使远程计算机屈服并开始发生超时(> 10 分钟)。

我们该如何解决这个问题呢?我希望能够做的是利用响应时间来限制调用。如果响应时间小于 30 秒,则继续添加线程,一旦超过 30 秒,就减少使用。有什么建议么?

注意与此问题中的响应相关:Parallel Foreach webservice call

c# parallel-processing multithreading task-parallel-library parallel.foreach

5
推荐指数
1
解决办法
568
查看次数

告诉 rsync 并行压缩

有没有办法告诉 rsync 并行进行压缩,例如。与pbzip2?我正在努力加快我们的夜间备份速度。我们的压缩 Rsync 目前的传输速率为

sent 1705628134 bytes  received 19432 bytes  3076010.04 bytes/sec
total size is 8064769536  speedup is 4.73
Run Code Online (Sandbox Code Playgroud)

而一个CPU核心开启99%,网络允许55-60Mbit/s,显然比实际的8*3Mbit/s要高。我们有 Oracle Linux 6.5

rsync.x86_64                           3.0.6-9.el6_4.1
Run Code Online (Sandbox Code Playgroud)

compression parallel-processing rsync

5
推荐指数
0
解决办法
1063
查看次数

使用maven进行分布式构建?

目前,我们有一个 Maven 项目,有几千个测试,需要 2 个小时才能运行。

我们尝试并行运行这些测试,但由于它们是功能测试,每个测试都以特定的方式配置系统,这会导致竞争条件和随机测试失败。

我想在 AWS 上启动 N 个服务器,然后让 Maven 将我的测试分开,并在这些服务器上运行它们(每个服务器将按顺序运行其测试,但所有服务器将并行运行),然后汇总结果。

有没有什么插件可以做这样的事情?

我已经看到了一些与我想要在 Jenkins 中实现的东西很接近的东西,但我更喜欢它是 Maven 驱动的,这样开发人员就可以在本地使用它,而无需安装 Jenkins。

parallel-processing distributed build maven

5
推荐指数
1
解决办法
1405
查看次数

Java并行流只使用一个线程?

我正在使用最新的 Java 8 lambda 和并行流来处理数据。我的代码如下:

ForkJoinPool forkJoinPool = new ForkJoinPool(10);
List<String> files = Arrays.asList(new String[]{"1.txt"}); 
List<String> result = forkJoinPool.submit(() ->
    files.stream().parallel()
        .flatMap(x -> stage1(x)) //at this stage we add more elements to the stream
        .map(x -> stage2(x))
        .map(x -> stage3(x))
        .collect(Collectors.toList())
).get();
Run Code Online (Sandbox Code Playgroud)

该流以一个元素开始,但在第二阶段添加更多元素。我的假设是该流应该并行运行,但在这种情况下仅使用一个工作线程。

如果我从 2 个元素开始(即,我将第二个元素添加到初始列表中),则会生成 2 个线程来处理流,依此类推...如果我没有显式地将流提交到 ForkJoinPool,也会发生这种情况。

问题是:它的行为是否记录在案,或者在实施过程中可能会发生变化?有什么方法可以控制这种行为并允许更多线程,无论初始列表如何?

java parallel-processing fork-join java-8 java-stream

5
推荐指数
1
解决办法
3554
查看次数