【问题标题】:Parallel *ply within functions函数内的并行 *ply
【发布时间】:2014-12-15 20:50:25
【问题描述】:

我想在函数中使用plyr 包的并行功能。

我原以为导出已在函数体中创建的对象(在此示例中,对象为df_2)的正确方法如下

# rm(list=ls())
library(plyr)
library(doParallel)

workers=makeCluster(2)
registerDoParallel(workers,core=2)

plyr_test=function() {
  df_1=data.frame(type=c("a","b"),x=1:2)
  df_2=data.frame(type=c("a","b"),x=3:4)

  #export df_2 via .paropts  
  ddply(df_1,"type",.parallel=TRUE,.paropts=list(.export="df_2"),.fun=function(y) {
    merge(y,df_2,all=FALSE,by="type")
  })
}
plyr_test()
stopCluster(workers)

但是,这会引发错误

Error in e$fun(obj, substitute(ex), parent.frame(), e$data) : 
  unable to find variable "df_2"

所以我做了一些研究,发现如果我手动导出df_2 就可以了

workers=makeCluster(2)
registerDoParallel(workers,core=2)

plyr_test_2=function() {
  df_1=data.frame(type=c("a","b"),x=1:2)
  df_2=data.frame(type=c("a","b"),x=3:4)

  #manually export df_2
  clusterExport(cl=workers,varlist=list("df_2"),envir=environment())

  ddply(df_1,"type",.parallel=TRUE,.fun=function(y) {
    merge(y,df_2,all=FALSE,by="type")
  })
}
plyr_test_2()
stopCluster(workers)

它给出了正确的结果

  type x.x x.y
1    a   1   3
2    b   2   4

但我也发现下面的代码有效

workers=makeCluster(2)
registerDoParallel(workers,core=2)

plyr_test_3=function() {
  df_1=data.frame(type=c("a","b"),x=1:2)
  df_2=data.frame(type=c("a","b"),x=3:4)

  #no export at all!
  ddply(df_1,"type",.parallel=TRUE,.fun=function(y) {
    merge(y,df_2,all=FALSE,by="type")
  })
}
plyr_test_3()
stopCluster(workers)

plyr_test_3() 也给出了正确的结果,我不明白为什么。我本来以为我必须导出df_2...

我的问题是:在函数内处理并行*ply 的正确方法是什么?显然,plyr_test() 是不正确的。我不知何故觉得plyr_test_2() 中的手动导出是没用的。但我也认为plyr_test_3() 是一种糟糕的编码风格。有人可以详细说明一下吗?谢谢大家!

【问题讨论】:

  • 附带说明:由于您使用的是ddply,您也可以尝试使用 dplyr,它是 plyr 的下一个版本对于数据帧,它可能会比ddply 的并行化更快地提高您的代码性能。见Introduction to dplyr
  • 谢谢。不知道那个。
  • 还有一个关于未来问题的注意事项:不要在您的问题中添加 rm(list=ls()) 未注释。其他人可能会在不注意的情况下运行代码,从而从他们的会话中删除重要数据。
  • :) 你是对的。对此感到抱歉。
  • 我上面给出的代码只是一个最小的例子。实际上,我什至没有在我的脚本中使用merge。

标签: r parallel-processing plyr


【解决方案1】:

plyr_test 的问题在于df_2 是在plyr_test 中定义的,而doParallel 包无法访问它,因此它在尝试导出df_2 时会失败。所以这是一个范围界定问题。 plyr_test2 避免了这个问题,因为它没有尝试使用.export 选项,但正如您所猜测的,不需要调用clusterExport。

plyr_test2 和plyr_test3 都成功的原因是df_2 与通过.fun 参数传递给ddply 函数的匿名函数一起序列化。事实上,df_1 和 df_2 都与匿名函数一起序列化,因为该函数是在 plyr_test2 和 plyr_test3 中定义的。在这种情况下包含df_2 很有帮助,但包含df_1 是不必要的,并且可能会影响您的性能。

只要在匿名函数的环境中捕获df_2,就不会使用df_2 的其他值,无论您导出什么。除非您可以阻止它被捕获,否则使用.export 或clusterExport 将其导出是没有意义的,因为将使用捕获的值。您只能通过尝试将其导出给工作人员来让自己陷入困境(就像您在.export 所做的那样)。

请注意,在这种情况下,foreach 不会自动导出 df_2,因为它无法分析匿名函数的主体以查看引用了哪些符号。如果您不使用匿名函数直接调用 foreach,那么它将看到引用并自动导出它,因此无需使用 .export 显式导出它。

您可以通过在将plyr_test 传递给ddply 之前修改其环境来防止plyr_test 的环境与匿名函数一起被序列化:

plyr_test=function() {
  df_1=data.frame(type=c("a","b"),x=1:2)
  df_2=data.frame(type=c("a","b"),x=3:4)
  clusterExport(cl=workers,varlist=list("df_2"),envir=environment())
  fun=function(y) merge(y, df_2, all=FALSE, by="type")
  environment(fun)=globalenv()
  ddply(df_1,"type",.parallel=TRUE,.fun=fun)
}

foreach 包的一个优点是它不鼓励您在可能意外捕获大量变量的另一个函数内创建函数。


这个问题向我表明foreach 应该包含一个名为.exportenv 的选项,该选项类似于clusterExport envir 选项。这对plyr 非常有帮助,因为它允许使用.export 正确导出df_2。但是,除非从 .fun 函数中删除包含 df_2 的环境,否则仍不会使用导出的值。

【讨论】:

  • 您如何解释 ARobertson 的带有虚拟变量 df_2='hi'; plyr_test("df_2",T); 的示例有效?显然,plyr_test("df_2",T) 之前的废话定义df_2='hi' 使得df_2 的导出工作。我不明白为什么会导出正确版本的df_2(在plyr_test 中定义)而不是df_2='hi'。
  • 我想我明白了。 plyr_test 中的df_2 无论plyr_test 之前是否有虚拟df_2='hi' 都会被导出。但是使用.export="df_2" 会触发ddply 在其搜索路径中搜索df_2。由于ddply 是在plyr_test 之外定义的函数,因此它只能在.GlobalEnv 中找到df_2。如果.GlobalEnv 中没有虚拟df_2='hi',它将引发错误。如果它找到df_2='hi',它也会尝试导出它。在后一种情况下,我想它会引发类似于我的foreach 示例中的警告。但该警告可能被忽略了。
  • @cryo111 关闭。 plyr_test 中的 df_2 被捕获,不管其他任何事情。除非在.GlobalEnv 中定义,否则使用.export="df_2" 将失败(尽管不会使用此导出值)。但是如果你使用.export="df_2",foreach 不会发出任何警告,因为当从ddply 调用时它不会自动导出.df_2,因为它无法分析匿名函数的主体。顺便说一句,我在回答中添加了更多解释来解决您的 cmets。
【解决方案2】:

看起来像是范围问题。

这是我的“测试套件”,它允许我 .export 不同的变量或避免在函数内创建 df_2。我在函数外部添加和删除一个虚拟 df_2 和 df_3 并进行比较。

library(plyr)
library(doParallel)

workers=makeCluster(2)
registerDoParallel(workers,core=2)

plyr_test=function(exportvar,makedf_2) {
  df_1=data.frame(type=c("a","b"),x=1:2)
  if(makedf_2){
    df_2=data.frame(type=c("a","b"),x=3:4)
  }
  print(ls())

  ddply(df_1,"type",.parallel=TRUE,.paropts=list(.export=exportvar,.verbose = TRUE),.fun=function(y) {
    z <- merge(y,df_2,all=FALSE,by="type")
  })
}
ls()
rm(df_2,df_3)
plyr_test("df_2",T)
plyr_test("df_2",F)
plyr_test("df_3",T)
plyr_test("df_3",F)
plyr_test(NULL,T) #ok
plyr_test(NULL,F)
df_2='hi'
ls()
plyr_test("df_2",T) #ok
plyr_test("df_2",F)
plyr_test("df_3",T)
plyr_test("df_3",F)
plyr_test(NULL,T) #ok
plyr_test(NULL,F)
df_3 = 'hi'
ls()
plyr_test("df_2",T) #ok
plyr_test("df_2",F)
plyr_test("df_3",T) #ok
plyr_test("df_3",F)
plyr_test(NULL,T) #ok
plyr_test(NULL,F)
rm(df_2)
ls()
plyr_test("df_2",T)
plyr_test("df_2",F)
plyr_test("df_3",T) #ok
plyr_test("df_3",F)
plyr_test(NULL,T) #ok
plyr_test(NULL,F)

我不知道为什么,但是 .export 在函数外部的全局环境中查找 df_2 (我在代码中看到了 parent.env(),这可能比全局环境“更正确”)而计算需要变量和ddply在同一个环境,并自动导出。

在函数外部使用 df_2 的虚拟变量允许 .export 工作,而计算使用内部的 df_2。

当.export在函数外找不到变量时,输出:

Error in e$fun(obj, substitute(ex), parent.frame(), e$data) : 
  unable to find variable "df_2" 

在函数外部有一个 df_2 虚拟变量但内部没有一个,.export 很好,但 ddply 输出:

Error in do.ply(i) : task 1 failed - "object 'df_2' not found"

这可能是因为这是一个小示例,或者可能无法并行化,它实际上是在一个内核上运行并且避免了导出任何内容的需要。如果没有 .export,更大的示例可能会失败,但其他人可以尝试。

【讨论】:

  • 有趣的测试套件!尤其是df_2 虚拟变量示例。请在下面查看我的“评论”(作为答案发布)。
  • 如果他们使用 parent.frame 代替,问题可能会得到解决的情况之一。这在函数内部也可以正常工作。我以前对这种范围界定问题感到头疼。 stackoverflow.com/questions/34395712/…
【解决方案3】:

感谢@ARobertson 的帮助! 非常有趣的是,plyr_test("df_2",T) 在函数体外部定义虚拟对象 df_2 时起作用。

看起来ddply 最终会调用llply,而后者又会调用foreach(...) %dopar% {...}。

我也尝试用foreach 重现该问题,但foreach 工作正常。

library(plyr)
library(doParallel)

workers=makeCluster(2)
registerDoParallel(workers,core=2)

foreach_test=function() {
  df_1=data.frame(type=c("a","b"),x=1:2)
  df_2=data.frame(type=c("a","b"),x=3:4)
  foreach(y=split(df_1,df_1$type),.combine="rbind",.export="df_2") %dopar% {
    #also print process ID to be sure that we really use different R script processes
    cbind(merge(y,df_2,all=FALSE,by="type"),Sys.getpid())
  }
}

foreach_test()
stopCluster(workers)

它会引发警告

Warning message:
In e$fun(obj, substitute(ex), parent.frame(), e$data) :
  already exporting variable(s): df_2

但它返回正确的结果

  type x.x x.y Sys.getpid()
1    a   1   3          216
2    b   2   4         1336

所以,foreach 似乎会自动导出df_2。事实上,foreachvignette 声明

... %dopar% 函数注意到这些变量被引用,并且 它们是在当前环境中定义的。在这种情况下 %dopar% 会自动将它们导出到并行执行 工人一次,并将它们用于所有表达式评估 那个 foreach 执行 ....

因此我们可以省略.export="df_2"而使用

library(plyr)
library(doParallel)

workers=makeCluster(2)
registerDoParallel(workers,core=2)

foreach_test_2=function() {
  df_1=data.frame(type=c("a","b"),x=1:2)
  df_2=data.frame(type=c("a","b"),x=3:4)
  foreach(y=split(df_1,df_1$type),.combine="rbind") %dopar% {
    #also print process ID to be sure that we really use different R script processes
    cbind(merge(y,df_2,all=FALSE,by="type"),Sys.getpid())
  }
}

foreach_test_2()
stopCluster(workers)

相反。这会在没有警告的情况下进行评估。

ARobertson 的虚拟变量示例以及 foreach 工作正常的事实让我现在认为 *ply 处理环境的方式存在问题。

我的结论是:

plyr_test_3() 和 foreach_test_2() 两个函数(不明确导出 df_2)运行时没有错误并给出相同的结果。因此,ddply 和 parallel=TRUE 基本上可以工作。但是使用更“详细”的编码风格(即显式导出df_2),例如plyr_test() 会引发错误,而foreach(...) %dopar% {...} 只会引发警告。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-02-25
    • 2022-11-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-24
    相关资源
    最近更新 更多