【问题标题】:Convert R apply statement to lapply for parallel processing将 R apply 语句转换为 lapply 以进行并行处理
【发布时间】:2017-10-02 19:45:10
【问题描述】:

我有以下 R “应用”语句:

for(i in 1:NROW(dataframe_stuff_that_needs_lookup_from_simulation))
{
    matrix_of_sums[,i]<-
    apply(simulation_results[,colnames(simulation_results) %in% 
    dataframe_stuff_that_needs_lookup_from_simulation[i,]],1,sum)
}

所以,我有以下数据结构:

simulation_results:具有列名的矩阵,用于标识 2000 次模拟(行)的每个可能的所需模拟查找数据。

dataframe_stuff_that_needs_lookup_from_simulation:除其他项目外,还包含其值与 Simulation_results 数据结构中的列名匹配的字段。

ma​​trix_of_sums: 函数运行时,一个 2000 行 x 250,000 列(模拟数量 x 正在模拟的项目)结构用于保存模拟结果。

因此,apply 函数在 250,000 个数据集中查找每一行的数据帧列值,计算总和,并将其存储在 matrix_of_sums 数据结构中。

很遗憾,此处理需要很长时间。我已经探索了使用 rowsums 作为替代方案,它已将处理时间缩短了一半,但我想尝试多核处理,看看是否可以进一步缩短处理时间。有人可以帮我将上面的代码从“apply”转换为“lapply”吗?

谢谢!

【问题讨论】:

  • applylapply 不会有太大的不同。但看起来你只是在做一个行求和,所以看看rowSums()。如果您想尝试并行计算,请尝试查看 foreach 包。
  • 为什么不与 rowsums 分享您的解决方案?
  • @MrFlick 嘿,感谢您的回复。实际上,我正在尝试转换为 lapply 以便我可以使用 mclapply 函数,该函数处理多个内核。
  • @Moody_Mudskipper rowsums 版本如下所示:matrix_of_sums[, i]

标签: r parallel-processing lapply


【解决方案1】:

使用base R并行,试试

library(parallel)
cl <- makeCluster(detectCores())
matrix_of_sums <- parLapply(cl, 1:nrow(dataframe_stuff_that_needs_lookup_from_simulation), function(i)
    rowSums(simulation_results[,colnames(simulation_results) %in% 
        dataframe_stuff_that_needs_lookup_from_simulation[i,]]))
stopCluster(cl)
ans <- Reduce("cbind", matrix_of_sums)

你也可以试试foreach %dopar%

library(doParallel)  # will load parallel, foreach, and iterators
cl <- makeCluster(detectCores())
registerDoParallel(cl)
matrix_of_sums <- foreach(i = 1:NROW(dataframe_stuff_that_needs_lookup_from_simulation)) %dopar% {
    rowSums(simulation_results[,colnames(simulation_results) %in% 
    dataframe_stuff_that_needs_lookup_from_simulation[i,]])
}
stopCluster(cl)
ans <- Reduce("cbind", matrix_of_sums)

我不太确定您最终希望输出的结果如何,但看起来您正在对每个结果进行cbind。但是,如果您期待别的东西,请告诉我。

【讨论】:

  • 您好,感谢您的回复。当我使用您的代码时,我收到以下错误:makePSOCKcluster(spec, ...) 中的错误:numeric 'names' must be >= 1. 我认为这可能是由于调用 detectCores() 引起的,它正在返回“NA”在我的机器上。我还不知道为什么。
  • 这很奇怪。如果您知道,您还可以手动指定可用内核的数量。以makeCluster(8) 为例。
  • 我想我们已经接近了。它正在检测节点的数量,但无法找到我正在操作的原始对象之一,simulation_results。这很奇怪,因为当我在控制台中输入它的名称时,我可以看到它充满了数据。我是否需要做任何其他事情才能将该对象传递给核心?这是我得到的错误: checkForRemoteErrors(val) 中的错误:8 个节点产生错误;第一个错误:找不到对象“simulation_results”
  • 好的,我添加了以下行:clusterExport(cl=cl, varlist=c("simulation_results")),按照这里的建议:stackoverflow.com/questions/12019638/… 我仍然遇到错误,但它得到了我已经过了最后一步。
  • 是的,重要的是导出您的数据 AND 包(如果您在 lapply 或 foreach 循环中使用任何包)。您可以使用基本 R 并行导出带有 clusterEvalQ(cl, { library(example)}) 的包,以及使用 foreach 导出 foreach(..., .packages=c("example")...) 的包
【解决方案2】:

实际上没有任何适用的或示例数据可供使用...该过程将如下所示:

  • 创建一个持有矩阵(matrix_of_sums)
  • 逐行循环遍历变量表(dataframe_stuff_that_needs_lookup_from_simulation)
  • 在仿真模型中查找匹配索引(simulation_results)
  • rowSums绑定到保持矩阵(总和矩阵)

我重新创建了一个无意义的样本集并产生相同的结果,但应该适用于您的数据

# Holding matrix which will be our end-goal
msums <- matrix(nrow = 2000,ncol = 0)
# Loop
    parallel::mclapply(1:nrow(ts_df), function(i){
       # Store the row to its own variable for ease
       d <- ts_df[i,]
       # cbind the results using the global assignment operator `<<-`
       msums <<- cbind(
                 msums, 
                 rowSums(
                    sim_df[,which(colnames(sim_df) %in% colnames(d))]
            ))
    }, mc.cores = parallel::detectCores(), mc.allow.recursive = TRUE)

【讨论】:

  • 嗨,Carl,接下来我将尝试您的解决方案。感谢您的回复!
  • Carl,这个版本也可以工作(我认为 - 我收到了一些核心错误,但工作完成了),但时间比运行 7 个核心的 CPak 稍差。您的版本将我的工作从 24 分钟缩短到 3'50"。非常感谢您的帮助!
  • 你运行的是什么环境?例如,我在具有 32 个内核的 ubuntu16.04 上使用 rstudio server 64x。 7 核听起来很奇怪,这就是我问的原因
  • 如果我们有实际数据也会很有帮助。因为 data.table 等可能会有更快的实现
  • 我在我的家庭桌面上运行 Rstudio 1.0.153 和 Ubuntu 16.04、一个八核的 Intel i7-6700K 和 32 GB 的 DDR4 RAM。最终,我想转移到 AWS。我在原始帖子中引用的数据框由许多与字符串字段相关联的浮点字段组成,这些字段需要与模拟矩阵的列匹配以进行重复查找。数据框中的浮点字段并不重要,可以删除;只有字符串字段以及它们与查找矩阵列名的匹配方式很重要。任何其他建议都会很棒!谢谢。
猜你喜欢
  • 1970-01-01
  • 2020-06-28
  • 1970-01-01
  • 2019-12-24
  • 2017-03-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多