【问题标题】:Running a user defined for-loop function in parallel on grouped data.table在分组的 data.table 上并行运行用户定义的 for 循环函数
【发布时间】:2020-01-15 17:14:59
【问题描述】:

我正在使用 R 中的 data.table 大约 6e6 行,我创建了一个函数,我通过 data.table 来创建一个基于两个分组值的新列。从技术上讲,我的函数循环遍历分组参数的每一行并执行一些非常简单的代数运算,但考虑到我的 data.table 的大小,这将需要相当长的时间。

我熟悉 foreach() 函数和其他使用多核进行计算的函数,但我还没有阅读或遇到过使用并行化来加速在函数中指定的 for 循环的方法通过data.table。本质上,我希望每个 for 循环迭代都由多个内核处理,而不是一个。有没有人在使用包含 for 循环的用户指定函数时有这方面的经验和/或在 data.table 中实现这一点?

【问题讨论】:

  • 请提供您尝试并行化的 for 循环的可重现示例。

标签: r for-loop parallel-processing data.table grouping


【解决方案1】:

我认为对您来说最好的答案可能是 Matt Dowle 和其他人正在齐心协力将 data.table 中的并行化内部化。老实说,我不能完全听懂所有的讨论,但我从经验中了解到分组现在在 data.table 中并行化,并且命令:

setDTthreads(0)

对我有帮助。这里有一些链接:

Is grouping parallelised in data.table 1.12.0? https://www.rdocumentation.org/packages/data.table/versions/1.12.2/topics/setDTthreads https://github.com/Rdatatable/data.table/issues/2031

我还想指出,dcast 处理分组非常快。我经常这样使用它:

# Some grouped Data:
z1[,1:4]
      CountyCode Gender  LYV Active
   1:         GY      M 2019   5742
   2:         KI      M 2019 244077
   3:         KI      F 2019 266944
   4:         CR      M 2018  51993
   5:         GY      M 2008    150
  ---                              
2172:         WT      U 2017      1
2173:         WK      M 2005      1
2174:         YA      U 1900     28
2175:         WL      U 1900      5
2176:         WK      U 1900      2

# Selecting group sums by gender for particular years:
z1[LYV %in% c("1900","2008","2012","2016","2017","2018","2019"),
.(Gender,LYV,Active)][,dcast(.SD,Gender ~ LYV,value.var="Active",fun.aggregate=sum)]

   Gender   1900  2008  2012   2016  2017   2018   2019
1:      F 275845 15694 43851 191024 27996 927968 777369
2:      M 307010 14543 41069 165942 24837 849066 688101
3:      U   6183    22    94   1161   233   5589   4804

最近对“dcast”的改进赋予了它很大的灵活性。在我读过的一些文档中,使用方括号将其传递给递归的 data.table ('.SD') 是不受欢迎的。但这对我来说效果很好。您也许可以简单地使用“set”命令和“dcast”来实现您的优化要求,而无需手动并行化。

【讨论】:

    【解决方案2】:

    由于您不提供示例数据, 这是一个简单的示例,可以帮助您入门。

    library(data.table)
    library(doParallel)
    
    dt <- data.table(a = sample(1:3, 1e6, TRUE),
                     b = sample(letters[1:5], 1e6, TRUE),
                     x = rnorm(1e6))
    
    workers <- makeCluster(detectCores())
    registerDoParallel(workers)
    
    ids <- dt[, .(list(.I)), by = .(a, b)]
    
    dt[unlist(ids$V1), y := foreach(i = ids$V1, .combine = c, .export = "dt", .packages = "data.table") %dopar% {
      setDT(dt)[i, as.numeric(scale(x))]
    }]
    
    stopCluster(workers); registerDoSEQ(); rm(workers)
    
    # sanity check
    dt[, identical(y, as.numeric(scale(x))), by = .(a, b)]
        a b   V1
     1: 2 c TRUE
     2: 1 a TRUE
     3: 3 d TRUE
     4: 1 d TRUE
     5: 1 b TRUE
     6: 3 c TRUE
     7: 2 e TRUE
     8: 3 e TRUE
     9: 2 a TRUE
    10: 2 b TRUE
    11: 2 d TRUE
    12: 1 c TRUE
    13: 3 a TRUE
    14: 3 b TRUE
    15: 1 e TRUE
    

    我们首先获取每个组的行索引并将它们保存在ids (在一个列表中,以便它们可以直接传递给foreach)。 分配y 的行将未列出的索引传递给data.tablei,以便将foreach 的结果分配给适当的行。

    我们在foreach 代码中使用setDT,因为表是序列化给工作人员的, 所以内存中的地址改变了 (至少我是这么认为的,也许其他人可以确认)。

    一定要用您的实际功能对其进行基准测试, 使用foreach 不能保证加速。 鉴于序列化, 数据的副本可能开销太大, 相对而言。

    【讨论】:

      猜你喜欢
      • 2018-07-29
      • 2023-02-07
      • 2021-05-10
      • 1970-01-01
      • 1970-01-01
      • 2016-12-29
      • 2023-03-14
      • 2017-04-18
      • 2021-10-12
      相关资源
      最近更新 更多