【问题标题】:Using foreach function to parallelise calculation使用 foreach 函数并行计算
【发布时间】:2019-02-18 05:05:03
【问题描述】:

我有一个包含 5000 个 csv 文件的文件夹,每个文件属于一个位置,包含从 1980 年到 2015 年的每日降雨量。文件的示例结构如下:

sample.file <- data.frame(location.id = rep(1001, times = 365 * 36), 
                      year = rep(1980:2015, each = 365),
                      day = rep(1:365, times = 36),
                      rainfall = sample(1:100, replace = T, 365 * 36))

我想读取一个文件并计算每年的总降雨量 并再次写入输出。有多种方法可以做到这一点:

方法一

for(i in seq_along(names.vec)){

  name <- namees.vec[i]
  dat <- fread(paste0(name,".csv"))

  dat <- dat %>% dplyr::group_by(year) %>% dplyr::summarise(tot.rainfall = sum(rainfall))

 fwrite(dat, paste0(name,".summary.csv"), row.names = F)
}

方法二:

my.files <- list.files(pattern = "*.csv")
dat <- lapply(my.files, fread)
dat <- rbindlist(dat)
dat.summary <- dat %>% dplyr::group_by(location.id, year) %>% 
               dplyr::summarise(tot.rainfall = sum(rainfall))

方法三:

我想使用foreach 来实现这一点。如何并行化上述任务 使用do parallelfor each 函数?

【问题讨论】:

  • 方法 4:fread 文件,rbind 它们并继续使用data.table 来提高性能(即allFilesBinded[, sum(rainfall), .(location.id, year)])怎么样?顺便说一句,因为1.11.0 fread is parallelized.
  • pbapply 包提供了简单的并行处理
  • 像 pblapply(my.files, fread, cl = mycl) 一样简单
  • 没有您的输入我无法测试,但我会选择:library(data.table); do.call(rbind, lapply(list.files(pattern = "*.csv"), fread))[, sum(rainfall), .(location.id, year)]
  • 通过this guide了解有关 {foreach} 并行性的更多信息。

标签: r foreach parallel-processing doparallel


【解决方案1】:

下面是您的foreach request 的骨架。

require(foreach)
require(doSNOW)
cl <- makeCluster(10, # number of cores, don't use all cores your computer have
                  type="SOCK") # SOCK for Windows, FORK for linux
registerDoSNOW(cl)
clusterExport(cl, c("toto", "truc"), envir=environment()) # R object needed for each core
clusterEvalQ(cl, library(tcltk)) # libraries needed for each core
my.files <- list.files(pattern = "*.csv")
foreach(i=icount(my.files), .combine=rbind, inorder=FALSE) %dopar% {
  # read csv file
  # estimate total rain
  # write output
}
stopCluster(cl)

但是当每次独立迭代的计算时间(CPU)高于其余操作时,并行化确实更好。在您的情况下,改进可能很低,因为每个内核都需要具有读取和写入的驱动器访问权限,并且由于写入是物理操作,因此最好按顺序进行(对硬件更安全,最终更高效与多个文件的共享位置相比,每个文件在驱动器中有独立的位置,需要索引等来区分它们以供您的操作系统使用——之前需要确认,这只是一个想法。

HTH

巴斯蒂安

【讨论】:

  • 并行化只应在读取文件时完成。模型估计(*)和保存到文件不应该是并行化的一部分。 (*) 是的,理论上你也许可以在这里做,因为它是一个总和。
  • 为了避免传达foreach() 是一个for循环而不是一个“应用”函数的风险,请考虑添加一个明确的返回值,例如可以dat.summary &lt;- foreach(...) 就像dat &lt;- lapply(my.files, fread) 一样。
【解决方案2】:

pbapply 包是最简单的并行方法

library (pbapply)

mycl <- makeCluster(4)
mylist <- pblapply(my.files, fread, cl = mycl)

【讨论】:

  • pbapply 不进行任何并行化。它所做的只是添加一个进度条。您可以将pblapply 替换为parallel::parLapply,它的工作原理完全相同。
  • 如果是关于添加进度条,这如何回答这个问题?
  • cl = mycl 启用并行。请尝试或阅读包参考。还有进度条
  • cl : 由 makeCluster 创建的集群对象,或一个整数来指示并行评估的子进程数(整数值在 Windows 上被忽略)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-09-08
  • 1970-01-01
  • 2019-03-30
  • 2014-11-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多