我家里有很多未使用过的电脑.对我来说,最简单的方法是利用它们并行化我的C#程序,只需很少或不需要更改代码?
我正在尝试的任务涉及循环许多英语句子,数据集可以很容易地分成更小的块,同时在不同的机器中处理.
我在 DataFrame 上进行了大型模拟df,我试图将模拟结果并行化并将模拟结果保存在名为 的 DataFrame 中simulation_results。
并行化循环工作得很好。问题是,如果我要将结果存储在数组中,我会将其声明为SharedArray循环之前。我不知道如何声明simulation_results为“共享数据帧”,它对所有处理器来说都可用并且可以修改。
代码片段如下:
addprocs(length(Sys.cpu_info()))
@everywhere begin
using <required packages>
df = CSV.read("/path/data.csv", DataFrame)
simulation_results = similar(df, 0) #I need to declare this as shared and modifiable by all processors
nsims = 100000
end
@sync @distributed for sim in 1:nsims
nsim_result = similar(df, 0)
<the code which for one simulation stores the results in nsim_result >
append!(simulation_results, nsim_result)
end
Run Code Online (Sandbox Code Playgroud)
问题在于,由于simulation_results未声明为由处理器共享和可修改,因此在循环运行后,它基本上会生成一个空的 DataFrame,如@everywhere simulation_results = similar(df, 0) …
使用我的语法使用简单的语法解析数百个文件
for @files -> $file {
my $input = $file.IO.slurp;
my $output = parse-and-convert($input);
$out-dir.IO.add($file ~ '.out').spurt: $output;
}
Run Code Online (Sandbox Code Playgroud)
循环相对较慢,在我的机器上大约需要 20 秒,因此我决定通过这样做来加快速度:
my @promises;
for @files -> $file {
my $input = $file.IO.slurp;
@promises.append: start parse-and-convert($input);
}
for await @promises -> $output {
$out-dir.IO.add($file ~ '.out').spurt: $output;
}
Run Code Online (Sandbox Code Playgroud)
这有效(至少在我的真实代码中,即对这个说明性示例中的任何拼写错误取模),但加速比我希望的要少得多:现在需要大约 11 秒,即我只获得了两倍。当然,这是值得重视的,但看起来存在很多争用,因为该程序使用的 CPU 数量少于 6 个(在有 16 个 CPU 的系统上),并且有相当多的开销(因为我没有得到6 加速都没有)。
我已经确认(通过插入一些say now - INIT.now)几乎所有运行时间都真正花费在内部await,如预期的那样,但我不知道如何进一步调试/分析它。我在 Linux 下执行此操作,因此我可以使用 perf,但我不确定它在 Raku 级别对我有何帮助。
有没有一些简单的方法可以提高这里的并行度?
编辑:为了说清楚,我可以忍受 20 秒(好吧,现在是 30 秒,因为我添加了更多东西)运行时间,我真的很好奇并行度是否可以在这里以某种方式提高,而不 …
我知道关于在 Julia 中使用 @threads、@distributed 和其他方法运行并行 for 循环存在很多问题。我尝试在那里实施解决方案,但没有成功。我想做的结构如下。
for index in list_of_indices
data = h5read("data_set_$index.h5")
result = perform_function(data)
save(result)
end
Run Code Online (Sandbox Code Playgroud)
数据集是独立的,并且该循环的任何部分都不依赖于任何其他部分。看来这应该是可并行的。
我尝试过,例如
“@threads for index in list_of_indices...”并且出现分段错误
“@distributed for index in list_of_indices...”并且代码实际上并未对我的数据执行该功能。
我想我错过了一些关于并行进程如何工作的信息,任何见解将不胜感激。
这是一个 MWE:
假设我们的工作目录中有文件 data_1.h5、data_2.h5、data_3.h5。(我不知道如何使事情比这更独立,因为我认为问题是由要求多个线程读取文件引起的。)
using Distributed
using HDF5
list = [1,2,3]
Threads.@threads for index in list
data = h5read("data_$index.h5", "data")
println(data)
end
Run Code Online (Sandbox Code Playgroud)
我得到的错误是
signal (11): Segmentation fault
signal (6): Aborted
Allocations: 1587194 (Pool: 1586780; Big: 414); GC: 1
Segmentation fault (core dumped)
Run Code Online (Sandbox Code Playgroud) 我对 ForEach-Object -Parallel 感到困惑。以下代码有包含超过 2000 个 blob 的 $blobs 数组。使用常规的foreach,我可以毫无问题地打印每个 blob 的名称。然后在第一个 foreach 之后使用ForEach-Object -Parallel ,不会打印任何内容。为什么 ?
foreach ($blob in $blobs) {
Write-Host $blob.Name
}
# Use parallel processing to process blobs concurrently
$blobs|ForEach-Object -Parallel {
param (
$blob)
Write-Host $blob.Name
} -ThrottleLimit 300
Run Code Online (Sandbox Code Playgroud) parallel-processing powershell azure-blob-storage foreach-object
我创建了各种参考类来适应一些 arima、garch 过程,并希望在并行计算中使用它们 parSapply
我先做了一些导出
cl <- makeCluster(mc <- getOption("cl.cores", 20))
clusterExport(cl, c("merge.xts", "index", "coredata", "xts", "lag.xts", "zoo", "LearnerPredict", "arima", "generic_learner", "arma_simple", "logwarn"))
clusterEvalQ(cl, "arma_simple")
clusterEvalQ(cl, "generic_learner")
generic_learner <- setRefClass(
Class = "generic_learner",
fields = list(
params = "list"
),
methods = list(
fitModel = function() {cat("overload function with fitting function \n")},
fcastModel = function() {cat("overload function with forecast function \n")},
fmt_params = function() {cat("overload function with formatted parameters \n")},
fmt_class = function() {cat("overload class\n")},
fmt_ref = function() {paste(.self$fmt_class(), .self$fmt_params(), …Run Code Online (Sandbox Code Playgroud) 我熟悉foreach,%dopar%之类的。我也是熟悉parallel的选项cv.glmnet。但是你如何设置嵌套的并行性如下?
library(glmnet)
library(foreach)
library(parallel)
library(doSNOW)
Npar <- 1000
Nobs <- 200
Xdat <- matrix(rnorm(Nobs * Npar), ncol = Npar)
Xclass <- rep(1:2, each = Nobs/2)
Ydat <- rnorm(Nobs)
Run Code Online (Sandbox Code Playgroud)
并行交叉验证:
cl <- makeCluster(8, type = "SOCK")
registerDoSNOW(cl)
system.time(mods <- foreach(x = 1:2, .packages = "glmnet") %dopar% {
idx <- Xclass == x
cv.glmnet(Xdat[idx,], Ydat[idx], nfolds = 4, parallel = TRUE)
})
stopCluster(cl)
Run Code Online (Sandbox Code Playgroud)
非并行交叉验证:
cl <- makeCluster(8, type = "SOCK")
registerDoSNOW(cl)
system.time(mods <- foreach(x …Run Code Online (Sandbox Code Playgroud) 很长一段时间以来,我一直在使用sfLapply来处理很多并行r脚本.然而,最近我已经深入研究并行计算,我一直在使用sfClusterApplyLB,如果单个实例不需要花费相同的时间来运行,那么可以节省大量时间.如果sfLapply将在加载新批处理之前等待批处理的每个实例完成(这可能导致空闲实例),完成任务的sfClusterApplyLB实例将立即分配给列表中的其余元素,因此可能会节省相当多的时间当实例没有花费相同的时间时.这让我质疑为什么我们在使用降雪时不想平衡我们的跑步?到目前为止我唯一发现的是,当并行脚本出现错误时,sfClusterApplyLB仍会在发出错误之前循环遍历整个列表,而sfLapply将在尝试第一批后停止.我还缺少什么?是否存在负载平衡的任何其他成本/缺点?下面是一个示例代码,显示了两者之间的差异
rm(list = ls()) #remove all past worksheet variables
working_dir="D:/temp/"
setwd(working_dir)
n_spp=16
spp_nmS=paste0("sp_",c(1:n_spp))
spp_nm=spp_nmS[1]
sp_parallel_run=function(sp_nm){
sink(file(paste0(working_dir,sp_nm,"_log.txt"), open="wt"))#######NEW
cat('\n', 'Started on ', date(), '\n')
ptm0 <- proc.time()
jnk=round(runif(1)*8000000) #this is just a redundant script that takes an arbitrary amount of time to run
jnk1=runif(jnk)
for (i in 1:length(jnk1)){
jnk1[i]=jnk[i]*runif(1)
}
ptm1=proc.time() - ptm0
jnk=as.numeric(ptm1[3])
cat('\n','It took ', jnk, "seconds to model", sp_nm)
#stop sinks
sink.reset <- function(){
for(i in seq_len(sink.number())){
sink(NULL)
}
}
sink.reset()
}
require(snowfall)
cpucores=as.integer(Sys.getenv('NUMBER_OF_PROCESSORS'))
sfInit( parallel=T, cpus=cpucores) # …Run Code Online (Sandbox Code Playgroud) 我的理解是 concurrent.futures 依靠酸洗参数来让它们在不同的进程(或线程)中运行。酸洗不应该创建参数的副本吗?在 Linux 上它似乎没有这样做,即,我必须明确地传递一个副本。
我试图理解以下结果:
<0> rands before submission: [17, 72, 97, 8, 32, 15, 63, 97, 57, 60]
<1> rands before submission: [97, 15, 97, 32, 60, 17, 57, 72, 8, 63]
<2> rands before submission: [15, 57, 63, 17, 97, 97, 8, 32, 60, 72]
<3> rands before submission: [32, 97, 63, 72, 17, 57, 97, 8, 15, 60]
in function 0 [97, 15, 97, 32, 60, 17, 57, 72, 8, 63]
in function 1 [97, 32, …Run Code Online (Sandbox Code Playgroud) python parallel-processing multiprocessing python-3.x concurrent.futures
我正在使用 OpenMP 并且我想生成线程,以便一个线程执行一段代码并完成,与运行并行 for 循环迭代的 N 个线程并行。
执行应该是这样的:
Section A (one thread) || Section B (parallel-for, multiple threads)
| || | | | | | | | | | |
| || | | | | | | | | | |
| || | | | | | | | | | |
| || | | | | | | | | | |
| || | | | | | | | | | |
V || …Run Code Online (Sandbox Code Playgroud) r ×3
julia ×2
.net ×1
c ×1
c# ×1
c++ ×1
cloud ×1
concurrency ×1
dataframe ×1
distributed ×1
foreach ×1
glmnet ×1
grammar ×1
nested ×1
openmp ×1
powershell ×1
python ×1
python-3.x ×1
raku ×1
snowfall ×1