我正在寻找一种“整洁”且有效的方法来实现长步骤 1(可以并行化)和步骤 2 的组合,步骤 2 需要按原始顺序(如果可能的话,尽量减少来自第一步保存在 RAM 中)同时允许第二步在第一个对象的步骤 1 中的数据可用时立即开始,并与步骤 2 一起提供更多数据。
为了更详细地说明这一点,我需要压缩大量图像(慢速 - 第 1 步),然后通过网络连接按顺序发送每个图像(第 2 步)。在任何阶段限制 RAM 中准备好的压缩数据块的数量也很重要,例如,如果发送 1000 张图像,我想将“已完成”但未发送的图像数量限制为(例如)线程数/使用的处理器。
我已经完成了这个的“手写”版本,使用了一组 Task 对象,但它看起来很混乱,而且我相信其他人一定有类似的需求,所以有没有更“标准”的方法来做到这一点? 理想情况下,我希望有 2 个代表的 Parallel.ForEach 变体 - 一个用于第 1 步,一个用于第 2 步,我希望标准覆盖之一(例如包含“localFinal”参数的覆盖)可能有所帮助,但在原来这些最后阶段是“每个线程”,而不是“每个委托”。
任何人都可以指出我现有的巧妙方法来实现这一目标吗?
我必须运行很多随机森林模型,所以我想在我的 8 核服务器上使用 doParallel 来加速这个过程。
然而,某些模型需要比其他模型更长的时间,甚至可能会引发错误。我想并行运行 8 个模型,如果模型抛出错误和/或被跳过,那么工作人员应该继续。每个模型结果都保存在硬盘上,以便我以后可以访问和组合它们。
TryCatch
Run Code Online (Sandbox Code Playgroud)
或者
.errorhandling="remove"
Run Code Online (Sandbox Code Playgroud)
没有解决问题。我得到
Error in unserialize(socklist[[n]]) : error reading from connection
Run Code Online (Sandbox Code Playgroud)
代码示例:我用 %do% 试了一下,模型 2-7 运行成功。然而在 %dopar% 我得到了显示的错误
foreach(model=1:8, .errorhandling="remove") %dopar% {
tryCatch({
outl <- rf_perform(...)
saveRDS(outl,file=getwd() %+% "/temp/result_" %+% model %+% ".rds")
}, error = function(e) {print(e)}, finally = {})
}
Run Code Online (Sandbox Code Playgroud) 我在 1000 多个数据集上cv.glmnet从glmnet包中并行运行。在每次运行中,我都会设置种子以使结果可重现。我注意到的是我的结果不同。问题是当我在同一天运行代码时,结果是一样的。但第二天他们就不同了。
这是我的代码:
model <- function(path, file, wyniki, faktor = 0.75) {
set.seed(2)
dane <- read.csv(file)
n <- nrow(dane)
podzial <- 1:floor(faktor*n)
########## GLMNET ############
nFolds <- 3
train_sparse <- dane[podzial,]
test_sparse <- dane[-podzial,]
# fit with cross-validation
tryCatch({
wart <- c(rep(0,6), "nie")
model <- cv.glmnet(train_sparse[,-1], train_sparse[,1], nfolds=nFolds, standardize=FALSE)
pred <- predict(model, test_sparse[,-1], type = "response",s=model$lambda.min)
# fetch of AUC value
aucp1 <- roc(test_sparse[,1],pred)$auc
}, error = function(e) print("error"))
results <- data.frame(auc = aucp1, n …Run Code Online (Sandbox Code Playgroud) 我正在尝试使用 MPI 写入文本文件,但未创建该文件。我只需要在主人处写(等级 = 0),但没有任何效果。它仅在我在控制台中运行程序(并保存损坏的元素)而不是在 Mpich2 中运行并且我附加了代码时才起作用。谢谢你的帮助。
/* -*- Mode: C; c-basic-offset:4 ; -*- */
/*
* (C) 2001 by Argonne National Laboratory.
* See COPYRIGHT in top-level directory.
*/
/* This is an interactive version of cpi */
#include <mpi.h>
#include <stdio.h>
#include <stdlib.h>
int main(int argc,char *argv[])
{
int namelen, numprocs, rank;
char processor_name[MPI_MAX_PROCESSOR_NAME];
MPI_Init(&argc,&argv);
MPI_Comm_rank(MPI_COMM_WORLD,&rank);
MPI_Comm_size(MPI_COMM_WORLD,&numprocs);
MPI_Get_processor_name(processor_name,&namelen);
MPI_Status status;
FILE* f = fopen("test.txt","wb+");
if (rank == 0) {
for (int i=0; i < 5; i++){
fprintf(f,"%d \n",i); …Run Code Online (Sandbox Code Playgroud) 我正在尝试运行如下所示的内容:
y = @parallel (min) for i in collection
f(i)
end
Run Code Online (Sandbox Code Playgroud)
wheref(i)是一个函数,它本质上是一个while循环,它计算满足其条件所需的迭代次数。开始时,终止条件之一是预定的迭代次数,n。但是,如果f(i)返回的值小于n理想n值,我想用 的值替换f(i)(例如,因为我正在寻找最小值f(i),如果f(j)是,m我希望所有其他循环停止检查它们是否达到m迭代)。
我是并行计算的新手,所以我可能会误解文档,但我认为我应该能够做这样的事情:
x = Channel{Int64}(1)
put!(x,n)
y = @parallel (min) for i in collection
f(i,x)
end
close(x)
Run Code Online (Sandbox Code Playgroud)
我已经修改f为采用Channel参数,现在它看起来像这样:
@everywhere function f(item,chan)
going = true
count = 0
while (going)
going = false
# perform some operations
if (count < fetch(chan) …Run Code Online (Sandbox Code Playgroud) 当我尝试运行 mpi 示例时,权限被拒绝。这是我尝试运行的代码。
#include <stdio.h>
#include <mpi.h>
int main (int argc,char *argv[])
{
int rank, size;
MPI_Init (&argc, &argv); /* starts MPI */
MPI_Comm_rank (MPI_COMM_WORLD, &rank); /* get current process id */
MPI_Comm_size (MPI_COMM_WORLD, &size); /* get number of processes */
printf( "Hello world from process %d of %d\n", rank, size );
MPI_Finalize();
return 0;
}
Run Code Online (Sandbox Code Playgroud)
我在主虚拟机上的共享文件夹中编译了它。我还生成了 ssh 密钥并将其复制到所有从属虚拟机。我有一个“主机”文件,其中包含所有虚拟机的所有 IP 地址,包括主虚拟机。
我用这个命令运行代码
`mpiexec -f hosts -n 4 hello_world
但我得到
===================================================================================
= BAD TERMINATION OF ONE OF YOUR APPLICATION PROCESSES
= …Run Code Online (Sandbox Code Playgroud) 我的luigi.cfg文件中有以下行(在所有节点、调度程序和工作程序上):
[core]
parallel-scheduling: true
Run Code Online (Sandbox Code Playgroud)
然而,当我在我的 luigi 调度程序上监控 CPU 利用率时(有大约 4000 个任务的图表,处理来自大约 100 个工作人员的请求),它只使用调度程序上的单个内核,luigid单线程经常达到 100% CPU 利用率. 我的理解是这个配置变量应该并行化任务的调度。
消息来源表明该标志确实应该在调度程序上使用多个内核。在https://github.com/spotify/luigi/blob/master/luigi/interface.py#L194 中,调用https://github.com/spotify/luigi/blob/master/luigi/worker。 py#L498.complete()并行检查任务的状态。
让我的 Luigi 调度程序利用其所有核心我还缺少什么?
我正在尝试了解 reduce 方法。如果我使用 reduce 和 stream() 我得到_ab,如果我使用 reduceparallelStream()我得到_a_b. 无论我们使用parallelStream还是stream,reduce的输出不应该是一样的吗?
import java.util.*;
import java.util.stream.*;
class TestParallelStream{
public static void main(String args[]){
List<String> l = Arrays.asList("a","b","c","d");
String join=l.stream()
.peek(TestParallelStream::sleepFor)
.reduce("_",(a,b) -> a.concat(b));
System.out.println(join);
}
public static void sleepFor(String w){
System.out.println("inside thread:"+w);
try{
Thread.currentThread().sleep(5000);
}catch(InterruptedException e){ }
}
}
Run Code Online (Sandbox Code Playgroud) 我正在分别具有 4 个和 8 个物理和逻辑内核的 PC(OS Linux)上运行以下代码(从doParallel 的 Vignettes 中提取)。
运行代码iter=1e+6或更少,一切都很好,我可以从 CPU 使用率中看到所有内核都用于此计算。然而,随着迭代次数的增多(例如iter=4e+6),在这种情况下并行计算似乎不起作用。当我还监视 CPU 使用率时,只有一个核心参与计算(100% 使用率)。
示例 1
require("doParallel")
require("foreach")
registerDoParallel(cores=8)
x <- iris[which(iris[,5] != "setosa"), c(1,5)]
iter=4e+6
ptime <- system.time({
r <- foreach(i=1:iter, .combine=rbind) %dopar% {
ind <- sample(100, 100, replace=TRUE)
result1 <- glm(x[ind,2]~x[ind,1], family=binomial(logit))
coefficients(result1)
}
})[3]
Run Code Online (Sandbox Code Playgroud)
你知道可能是什么原因吗?记忆可能是原因吗?
我四处搜索,发现这与我的问题有关,但重点是我没有出现任何错误,而且 OP 似乎通过在内部提供必要的包来提出解决方案foreach循环。但是可以看出,我的循环中没有使用任何包。
更新1
我的问题还是没有解决。根据我的实验,我不认为记忆可能是原因。我在运行以下简单并行(在所有 8 个逻辑内核上)迭代的系统上有 8GB 内存:
例2
require("doParallel")
require("foreach")
registerDoParallel(cores=8)
iter=4e+6
ptime <- system.time({
r <- foreach(i=1:iter, …Run Code Online (Sandbox Code Playgroud) 在 Julia 中,我想在模块内部定义的函数中使用addprocs和pmap。这是一个愚蠢的例子:
module test
using Distributions
export g, f
function g(a, b)
a + rand(Normal(0, b))
end
function f(A, b)
close = false
if length(procs()) == 1 # If there are already extra workers,
addprocs() # use them, otherwise, create your own.
close = true
end
W = pmap(x -> g(x, b), A)
if close == true
rmprocs(workers()) # Remove the workers you created.
end
return W
end
end
test.f(randn(5), 1)
Run Code Online (Sandbox Code Playgroud)
这将返回一个很长的错误
WARNING: Module test …Run Code Online (Sandbox Code Playgroud) r ×3
c ×2
doparallel ×2
julia ×2
c# ×1
glmnet ×1
java ×1
java-stream ×1
linux ×1
luigi ×1
mpi ×1
python ×1
random-seed ×1
reduce ×1
text-files ×1