【问题标题】:doParallel (package) foreach does not work for big iterations in RdoParallel (包) foreach 不适用于 R 中的大迭代
【发布时间】:2016-10-11 13:56:27
【问题描述】:

我在分别具有 4 个和 8 个物理内核和逻辑内核的 PC (OS Linux) 上运行以下代码(从 doParallel's 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]

您知道可能是什么原因吗?记忆可能是原因吗?

我四处搜索,发现 THIS 与我的问题相关,但关键是我没有遇到任何错误,并且 OP 似乎通过在 foreach 循环中提供必要的包来提出解决方案。但可以看出,我的循环中没有使用任何包。

更新1

我的问题还没有解决。根据我的实验,我不认为记忆可能是原因。我在运行以下简单并行(在所有 8 个逻辑内核上)迭代的系统上有 8GB 内存:

示例2

require("doParallel")
require("foreach")

registerDoParallel(cores=8)
iter=4e+6
ptime <- system.time({
    r <- foreach(i=1:iter, .combine=rbind) %dopar% {
        i
    }
})[3]

我运行这段代码没有问题,但是当我监控 CPU 使用率时,只有一个核心(共 8 个)是 100%。

更新2

至于 Example2,@SteveWeston(感谢您指出这一点)表示(在 cmets 中):“您更新中的示例遇到了小任务。只有大师有任何实际工作to do,包括发送任务和处理结果。这与原始示例的问题根本不同,原始示例确实在较少次数的迭代中使用了多个内核。"

但是,Example1 仍未解决。当我运行它并使用htop 监控进程时,更详细的情况如下:

让我们将所有 8 个创建的进程命名为 p1p8p1 的状态(htop 中的 S 列)是 R,这意味着它正在运行并且保持不变。然而,对于p2p8,几分钟后,状态变为D(即不间断睡眠),几分钟后,再次变为Z(即终止但未被其父级接收) .你知道为什么会这样吗?

【问题讨论】:

  • 您是否尝试过显式创建和注册集群?例如cl &lt;- makePSOCKcluster(8); registerDoParallel(cl)?
  • 即使iter=15e+6?您能否评论一下您的硬件规格(CPU 和内存)以及运行代码的操作系统类型?
  • 是的,它可能确实与内存或计算资源有关。那么你真的需要那么多迭代吗?如果是这样,您是否可以通过将其分成具有较少迭代的步骤然后组合(可能是平均)结果来获得相同的结果?我不确定您的用例是什么,但对于机器学习,您通常可以将模型进度保存到存储内存(即检查点)中,这样您就可以将大型神经网络和其他模型的训练分成多个会话。
  • 我让它在 Windows 上运行全部 8 个(在 10-15 秒后自行排序),I7 4910 逐字运行您的上述代码。
  • 您的更新中的示例遇到了小任务。只有master有真正的工作要做,包括发送任务和处理结果。这与在少量迭代中使用多个内核的原始示例的问题根本不同。

标签: r parallel-processing parallel-foreach doparallel


【解决方案1】:

起初我以为您遇到了内存问题,因为提交许多任务确实会使用更多内存,这最终会导致主进程陷入困境,因此我的原始答案显示了几种使用更少内存的技术。但是,现在听起来好像有一个启动和关闭阶段,只有主进程在忙,而工作人员在中间忙了一段时间。我认为问题在于此示例中的任务并不是真正的计算密集型任务,因此当您有很多任务时,您会开始真正注意到启动和关闭时间。我对实际计算进行了计时,发现每个任务只需要大约 3 毫秒。过去,你不会从并行计算中获得任何好处,任务那么小,但现在,取决于你的机器,你可以获得一些好处,但开销很大,所以当你有很多任务时,你真的会注意到开销。

我仍然认为我的其他答案对这个问题很有效,但是由于你有足够的内存,所以它是矫枉过正的。使用分块的最重要技术。这是一个使用分块的示例,对原始示例的更改很小:

require("doParallel")
nw <- 8
registerDoParallel(nw)
x <- iris[which(iris[,5] != "setosa"), c(1,5)]
niter <- 4e+6
r <- foreach(n=idiv(niter, chunks=nw), .combine='rbind') %dopar% {
  do.call('rbind', lapply(seq_len(n), function(i) {
    ind <- sample(100, 100, replace=TRUE)
    result1 <- glm(x[ind,2]~x[ind,1], family=binomial(logit))
    coefficients(result1)
  }))
}

请注意,这与我的其他答案略有不同。通过使用 idiv chunks 选项,而不是 chunkSize 选项,它只使用每个工作人员一个任务。这减少了 master 完成的工作量,如果你有足够的内存,这是一个很好的策略。

【讨论】:

  • 谢谢,投了赞成票。我会尽快完成这个问题。
  • 如果Example1 中的并行方法由于任务不是计算密集型而不起作用,那么为什么它适用于(我的意思是并行)小迭代而不是大迭代?
【解决方案2】:

我认为您的内存不足。这是该示例的修改版本,当您有很多任务时应该会更好。它使用 doSNOW 而不是 doParallel,因为 doSNOW 允许您在工作人员返回时使用 combine 函数处理结果。此示例将这些结果写入文件以使用更少的内存,但是它在最后使用“.final”函数将结果读回内存,但如果您没有足够的内存,您可以跳过它。

library(doSNOW)
library(tcltk)
nw <- 4  # number of workers
cl <- makeSOCKcluster(nw)
registerDoSNOW(cl)

x <- iris[which(iris[,5] != 'setosa'), c(1,5)]
niter <- 15e+6
chunksize <- 4000  # may require tuning for your machine
maxcomb <- nw + 1  # this count includes fobj argument
totaltasks <- ceiling(niter / chunksize)

comb <- function(fobj, ...) {
  for(r in list(...))
    writeBin(r, fobj)
  fobj
}

final <- function(fobj) {
  close(fobj)
  t(matrix(readBin('temp.bin', what='double', n=niter*2), nrow=2))
}

mkprogress <- function(total) {
  pb <- tkProgressBar(max=total,
                      label=sprintf('total tasks: %d', total))
  function(n, tag) {
    setTkProgressBar(pb, n,
      label=sprintf('last completed task: %d of %d', tag, total))
  }
}
opts <- list(progress=mkprogress(totaltasks))
resultFile <- file('temp.bin', open='wb')

r <-
  foreach(n=idiv(niter, chunkSize=chunksize), .combine='comb',
          .maxcombine=maxcomb, .init=resultFile, .final=final,
          .inorder=FALSE, .options.snow=opts) %dopar% {
    do.call('c', lapply(seq_len(n), function(i) {
      ind <- sample(100, 100, replace=TRUE)
      result1 <- glm(x[ind,2]~x[ind,1], family=binomial(logit))
      coefficients(result1)
    }))
  }

我包含了一个进度条,因为这个示例需要几个小时才能执行。

请注意,此示例还使用 iterators 包中的 idiv 函数来增加每个任务的工作量。这种技术称为chunking,通常可以提高并行性能。但是,使用idiv 会弄乱任务索引,因为变量i 现在是每个任务的索引而不是全局索引。对于全局索引,你可以编写一个自定义迭代器来包装idiv

idivix <- function(n, chunkSize) {
  i <- 1
  it <- idiv(n, chunkSize=chunkSize)
  nextEl <- function() {
    m <- nextElem(it)  # may throw 'StopIterator'
    value <- list(i=i, m=m)
    i <<- i + m
    value
  }
  obj <- list(nextElem=nextEl)
  class(obj) <- c('abstractiter', 'iter')
  obj
}

此迭代器发出的值是列表,每个列表都包含一个起始索引和一个计数。这是一个使用此自定义迭代器的简单 foreach 循环:

r <- 
  foreach(a=idivix(10, chunkSize=3), .combine='c') %dopar% {
    do.call('c', lapply(seq(a$i, length.out=a$m), function(i) {
      i
    }))
  }

当然,如果任务的计算量足够大,您可能不需要分块,可以使用原始示例中的简单 foreach 循环。

【讨论】:

  • 感谢您的解决方案,但是当我运行您的代码时,我得到:Warning: progress function failed: object 'pb' not found
  • @m0h3n 您始终可以通过不使用 foreach .options.snow 选项来禁用进度条。 pbprogress 是否都在全局环境中定义?
  • 是的,我只是复制/粘贴您的代码并运行它。但是看看comb 函数及其参数和for 里面的循环。这正常吗?
  • @m0h3n comb 函数使用第一个参数与后续参数不同,后者需要使用 foreach .init 参数。我已经将这种模式的示例作为示例包含在一些与 foreach 相关的包中。
  • 我想知道您是否在发布之前运行您的代码。我明白了:Error: object 'opts' not found。请尝试清理您的环境并运行您在此处发布的代码。
猜你喜欢
  • 1970-01-01
  • 2018-05-18
  • 2020-11-18
  • 2023-03-12
  • 2013-09-12
  • 2016-10-08
  • 1970-01-01
  • 2020-07-13
  • 2020-06-07
相关资源
最近更新 更多