【问题标题】:Memory issue with foreach loop in R on Windows 8 (64-bit) (doParallel package)Windows 8(64 位)上 R 中的 foreach 循环的内存问题(doParallel 包)
【发布时间】:2014-01-15 14:32:11
【问题描述】:

我正在尝试从串行方法转变为并行方法,以在大型 data.table 上完成一些多元时间序列分析任务。该表包含许多不同组的数据,我正在尝试使用doParallel 包从for 循环移动到foreach 循环,以利用安装的多核处理器。

我遇到的问题与内存以及新的 R 进程似乎如何消耗大量内存有关。我认为正在发生的事情是包含所有数据的大型 data.table 被复制到每个新进程中,因此我用完了 RAM,Windows 开始交换到磁盘。

我创建了一个简化的可重现示例,它复制了我的问题,但循环内的数据和分析更少。如果存在一种解决方案,该解决方案只能按需将数据外包给工作进程,或者在内核之间共享已经使用的内存,那将是理想的。或者,可能已经存在某种解决方案将大数据分成 4 个块并将它们传递给内核,以便它们有一个子集可以使用。

here on Stackoverflow 之前已经发布了一个类似的问题,但是我无法使用提供的bigmemory 解决方案,因为我的数据包含一个字符字段。我将进一步研究iterators 包,但我会感谢有此问题实践经验的成员提供的任何建议。

rm(list=ls())
library(data.table)
num.series = 40    # can customise the size of the problem (x10 eats my RAM)
num.periods = 200  # can customise the size of the problem (x10 eats my RAM)
dt.all = data.table(
    grp     = rep(1:num.series,each=num.periods), 
    pd      = rep(1:num.periods, num.series), 
    y       = rnorm(num.series * num.periods),
    x1      = rnorm(num.series * num.periods),
    x2      = rnorm(num.series * num.periods)
) 
dt.all[,y_lag := c(NA, head(y, -1)), by = c("grp")]

f_lm = function(dt.sub, grp) {
    my.model = lm("y ~ y_lag + x1 + x2 ", data = dt.sub)
    coef = summary(my.model)$coefficients
    data.table(grp, variable = rownames(coef), coef)
}

library(doParallel)  
registerDoParallel(4) 
foreach(grp=unique(dt.all$grp), .packages="data.table", .combine="rbind")  %dopar%
{
    dt.sub = dt.all[grp == grp]
    f_lm(dt.sub, grp)
}
detach(package:doParallel)

【问题讨论】:

  • 语句 dt.sub = dt.all[grp == grp] 将导致 dt.sub 的大量副本四处浮动。您是否尝试过在 lm 调用中使用 lm 的子集 = 参数来避免这样做?另外我想知道 f_lm 的返回值在计算时做了什么,它没有分配给任何东西。
  • 是的,这个例子有缺陷。实际上,我必须缩小stepAIC 的上限和下限的可能公式,过滤掉NAs 的行,获得最佳模型,记录系数,计算弹性和样本内精度等。所以有除了运行lm 之外还有几个步骤,我会尴尬地将data.tables 的列表返回给foreach。我一定会给subset一些考虑,谢谢你的建议。我可以处理dt.sub 的副本,只是dt.all 在内存中不能有很多副本。
  • 我很惊讶会创建 dt.all 的多个副本。根据多核的文档,“现代操作系统使用写时复制方法,这使得这对并行计算非常有吸引力,因为只有在计算期间修改的对象才会被实际复制,而所有其他内存都是直接共享的。”我知道 dt.all 在 Linux 上不会多次驻留在内存中。不确定 Windows。
  • doParallel: "多核功能仅在那些支持 fork 系统调用的操作系统上支持多个工作线程;这不包括 Windows。默认情况下,doParallel 在类 Unix 系统上使用多核功能,在 Windows 上使用雪功能”。我没有关于环境中的数据是否被复制的书面证据,但鉴于必须指定所需的包,并且我假设加载到 R 的每个实例中,那么数据也是必需的。我猜这是 Windows 的问题。
  • 是的。有人告诉我这个功能在 Unix 上效果更好(因为 fork 命令),但切换可能不是你的选择。

标签: r memory foreach data.table


【解决方案1】:

答案需要iterators 包和isplit 的使用,这与split 类似,因为它将主数据对象基于一个或多个factor 列分成块。 foreach 循环遍历数据块,仅将子集传递给工作进程,而不是整个表。

所以代码的区别如下:

library(iterators)
dt.all = data.table(
    grp     = factor(rep(1:num.series, each  =num.periods)),  # grp column is a factor
    pd      = rep(1:num.periods, num.series), 
    y       = rnorm(num.series * num.periods),
    x1      = rnorm(num.series * num.periods),
    x2      = rnorm(num.series * num.periods)
) 

results = 
    foreach(dt.sub = isplit(dt.all, dt.all$grp), .packages="data.table", .combine="rbind")  
    %dopar%
    {
        f_lm(dt.sub$value, dt.sub$key[[1]])
    }

isplit 的结果是 dt.sub 现在是具有 2 个元素的 listkey 本身是用于拆分的值列表,value 包含子集作为data.table.

此解决方案的功劳归功于SO answer given by David,Russell 对我在excellent blog post about iterators 上的问题作出了回应。

------------------------------------ 编辑 ---------- --------------------------

为了测试isplitDT v isplitrbindlist v rbind 的性能,使用了以下代码:

rm(list=ls())
library(data.table) ; library(iterators)  ;   library(doParallel)
num.series = 400
num.periods = 2000
dt.all = data.table(
    grp     = factor(rep(1:num.series,each=num.periods)), 
    pd      = rep(1:num.periods, num.series), 
    y       = rnorm(num.series * num.periods),
    x1      = rnorm(num.series * num.periods),
    x2      = rnorm(num.series * num.periods)
)
dt.all[,y_lag := c(NA, head(y, -1)), by = c("grp")]

f_lm = function(dt.sub, grp) {
    my.model = lm("y ~ y_lag + x1 + x2 ", data = dt.sub)
    coef = summary(my.model)$coefficients
    data.table(grp, variable = rownames(coef), coef)
}

registerDoParallel(8)

isplitDT <- function(x, colname, vals) {
  colname <- as.name(colname)
  ival <- iter(vals)
  nextEl <- function() {
    val <- nextElem(ival)
    list(value=eval(bquote(x[.(colname) == .(val)])), key=val)
  }
  obj <- list(nextElem=nextEl)
  class(obj) <- c('abstractiter', 'iter')
  obj
}

dtcomb <- function(...) {
  rbindlist(list(...))
}

# isplit/rbind
st1 = system.time(results <- foreach(dt.sub=isplit(dt.all,dt.all$grp),   
                    .combine="rbind",
                    .packages="data.table")  %dopar% {
    f_lm(dt.sub$value, dt.sub$key[[1]])
})
# isplit/rbindlist
st2 = system.time(results <- foreach(dt.sub=isplit(dt.all,dt.all$grp),  
                .combine='dtcomb', .multicombine=TRUE,
                .packages="data.table") %dopar% {
    f_lm(dt.sub$value, dt.sub$key[[1]])
})
# isplitDT/rbind
st3 = system.time(results <- foreach(dt.sub=isplitDT(dt.all, 'grp',     unique(dt.all$grp)),
            .combine='dtcomb', .multicombine=TRUE,
            .packages='data.table') %dopar% {
    f_lm(dt.sub$value, dt.sub$key)
})
# isplitDT/rbindlist
st4 = system.time(results <- foreach(dt.sub=isplitDT(dt.all, 'grp', unique(dt.all$grp)),
                .combine='dtcomb', .multicombine=TRUE,
                .packages='data.table') %dopar% {
    f_lm(dt.sub$value, dt.sub$key)
  })

rbind(st1, st2, st3, st4)

这给出了以下时间:

    user.self sys.self elapsed user.child sys.child
st1     12.08     1.53   14.66         NA        NA
st2     12.05     1.41   14.08         NA        NA
st3     45.33     2.40   48.14         NA        NA
st4     45.00     3.30   48.70         NA        NA

------------------------------------ 编辑 2 --------- --------------------------

感谢 Steve 的更新答案和函数 isplitDT2,它使用了 data.table 上的键,我们在速度方面有了一个明显的新赢家。运行microbenchmark 来比较我的原始解决方案(在这个答案中)显示从isplitDT2rbindlist 提高了大约7 倍。内存使用情况还没有直接比较,但是性能提升让我终于接受了答案。

【讨论】:

  • 我更新的答案中的 isplitDB2 在速度或内存使用方面做得更好吗?
  • @SteveWeston:感谢更新,它在问题中的玩具示例上表现得更好。我现在将在真实数据集上尝试它并评估内存使用情况。
【解决方案2】:

迭代器有助于减少需要传递给并行程序工作人员的内存量。由于您使用的是 data.table 包,因此最好使用迭代器并组合针对 data.table 对象优化的函数。例如,这里有一个类似 isplit 的函数,它适用于 data.table 对象:

isplitDT <- function(x, colname, vals) {
  colname <- as.name(colname)
  ival <- iter(vals)
  nextEl <- function() {
    val <- nextElem(ival)
    list(value=eval(bquote(x[.(colname) == .(val)])), key=val)
  }
  obj <- list(nextElem=nextEl)
  class(obj) <- c('abstractiter', 'iter')
  obj
}

请注意,它与isplit 不完全兼容,因为参数和返回值略有不同。可能还有更好的方法来对data.table进行子集化,但我认为这比使用isplit更有效。

这是您使用isplitDT 和使用rbindlist 的组合函数的示例,它比rbind 更快地组合data.tables:

dtcomb <- function(...) {
  rbindlist(list(...))
}

results <- 
  foreach(dt.sub=isplitDT(dt.all, 'grp', unique(dt.all$grp)),
          .combine='dtcomb', .multicombine=TRUE,
          .packages='data.table') %dopar% {
    f_lm(dt.sub$value, dt.sub$key)
  }

更新

我编写了一个名为isplitDT2 的新迭代器函数,它的性能比isplitDT 好得多,但要求data.table 有一个键:

isplitDT2 <- function(x, vals) {
  ival <- iter(vals)
  nextEl <- function() {
    val <- nextElem(ival)
    list(value=x[val], key=val)
  }
  obj <- list(nextElem=nextEl)
  class(obj) <- c('abstractiter', 'iter')
  obj
}

这被称为:

setkey(dt.all, grp)
results <-
  foreach(dt.sub=isplitDT2(dt.all, levels(dt.all$grp)),
          .combine='dtcomb', .multicombine=TRUE,
          .packages='data.table') %dopar% {
    f_lm(dt.sub$value, dt.sub$key)
  }

这使用二分搜索来子集dt.all,而不是矢量扫描,因此效率更高。但是,我不知道为什么isplitDT 会使用更多内存。由于您使用的是doParallel,它不会在发送任务时即时调用迭代器,您可能需要尝试拆分dt.all,然后将其删除以减少内存使用:

dt.split <- as.list(isplitDT2(dt.all, levels(dt.all$grp)))
rm(dt.all)
gc()
results <- 
  foreach(dt.sub=dt.split,
          .combine='dtcomb', .multicombine=TRUE,
          .packages='data.table') %dopar% {
    f_lm(dt.sub$value, dt.sub$key)
  }

这可能有助于减少主进程在执行 foreach 循环期间所需的内存量,同时仍然只将所需的数据发送给工作人员。如果您仍然有内存问题,您还可以尝试使用 doMPI 或 doRedis,两者都根据需要获取迭代器值,而不是一次全部获取,从而提高内存效率。

【讨论】:

  • 很好的解决方案史蒂夫,特别是解决方案的组合功能和通用性,谢谢。不幸的是,我不能接受我自己的答案,因为这个显然更好!
  • 我尝试在我的真实数据上使用这种方法,但它导致笔记本电脑内存不足。在具有 10GB RAM 的集群上也存在同样的问题。因此,我根据原始问题规范测试了这两种解决方案,并惊讶地发现system.timeisplitDTdtcomb 的比较为49 秒,而基本isplitrbind 的比较为14.76 秒。所以我进行了 4 次测试来评估rbind/rbindlist x isplit/isplitDT 的组合。这表明放缓是由于isplitDT 而不是rbindlist。如果需要,将发布测试代码。有什么想法吗?
【解决方案3】:

将所有内容保存在内存中是R 程序员必须学会处理的那些(啊,烦人的)事情之一。很容易将您的代码示例想象为受内存限制或受 CPU 限制,您需要在尝试应用变通方法之前弄清楚这一点。

假设内存被您的数据集 (dt_all) 消耗,而不是在实际模型运行期间,您可能能够释放足够的内存供工作进程并行化:

foreach(grp=unique(dt.all$grp), .packages="data.table", .combine="rbind")  %dopar%
{
    dt.sub = dt.all[grp == grp]
    rm(dt.all)
    gc()
    f_lm(dt.sub, grp)
}

但是,这假设您的工作集 (dt.sub) 足够小,以至于您一次可以在内存中容纳多个。不难想象一个问题集太大了。此外,这真的很烦人,所有的工作人员都会同时启动并杀死你的机器,所以你可能需要让他们暂停几秒钟,让其他孩子加载并释放内存。

尽管极其愚蠢和蛮力,我已经通过将子集作为单独的数据文件写入磁盘来处理这个确切的问题,然后使用批处理脚本并行运行我的计算。

【讨论】:

  • 关于删除dt.all 和子负载的阶段性非常好!幸运的是,我的dt.sub 并没有太大(313 x 100)以至于消耗大量内存。正如您所建议的那样,写入磁盘将是一种选择,我明天将对其进行测试。目前我正在笔记本电脑上玩并行计算,在生产中我将能够使用大学 HPC,所以当我了解我在做什么时应该能够使用更多的内存和节点!
  • 如果您的工作集是可管理的,那么您的主要问题可能是在创建子集的时间和运行模型的时间之间进行权衡。这种权衡可能会决定你的方法。 GL!
  • 如果任务的数量大于工作人员的数量,删除 dt.all 将导致错误,因为当任何工作人员执行后续任务时,任务将开始失败。当使用默认的 doSEQ 后端时,这是一个更大的问题,因为在执行第一个任务后会删除主副本。
猜你喜欢
  • 2017-10-19
  • 2016-05-09
  • 2020-07-13
  • 2020-06-07
  • 1970-01-01
  • 1970-01-01
  • 2021-09-16
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多