【问题标题】:How to implement lapply function in R using package "sparklyr"如何使用包“sparklyr”在 R 中实现 lapply 功能
【发布时间】:2018-06-24 04:51:27
【问题描述】:

我对 Spark 很陌生,我尝试在网上寻找一些东西,但我没有找到任何令人满意的东西。

我一直使用命令mclapply 运行并行计算,我喜欢它的结构(即,第一个参数用作滚动索引,第二个参数是要并行化的函数,然后是传递给函数的其他可选参数)。 现在我正在尝试通过 Spark 做同样的事情,即我想在 Spark 集群的所有节点之间分配我的计算。这就是我所学到的以及我认为代码应该如何构造的内容(我正在使用包sparklyr):

  1. 我使用命令 spark_connect 创建到 Spark 的连接;
  2. 我使用copy_to 在 Spark 环境中复制我的 data.frame,并通过其 tibble 访问它;
  3. 我想实现mclapply的“Spark-friendly”版本,但我看到包中没有类似的功能(我看到SparkR包中存在函数spark.lapply,但不幸的是,它不再在 CRAN 中了)。

下面是我实现的一个简单的测试脚本,它使用函数mclapply 工作。

#### Standard code that works with mclapply #########
dfTest = data.frame(X = rep(1, 10000), Y = rep(2, 10000))

.testFunc = function(X = 1, df, str) {
    rowSelected = df[X, ]
    y = as.numeric(rowSelected[1] + rowSelected[2])
    return(list(y = y, str = str))
}

lOutput = mclapply(X = 1 : nrow(dfTest), FUN = .testFunc, df = dfTest, 
                   str = "useless string", mc.cores = 2)

######################################################

###### Similar code that should work with Spark ######
library(sparklyr)
sc = spark_connect(master = "local")

dfTest = data.frame(X = rep(1, 10000), Y = rep(2, 10000))

.testFunc = function(X = 1, df, str) {
  rowSelected = df[X, ]
  nSum = as.numeric(rowSelected[1] + rowSelected[2])
  return(list(nSum = nSum, str = str))
}

dfTest_tbl = copy_to(sc, dfTest, "test_tbl", overwrite = TRUE)

# Apply similar function mclapply to dfTest_tbl, that works with 
# Spark
# ???
######################################################

如果有人已经找到了解决方案,那就太好了。其他参考/指南/链接也非常受欢迎。谢谢!

【问题讨论】:

    标签: r apache-spark parallel-processing sparklyr mclapply


    【解决方案1】:

    sparklyr

    spark_apply 是您正在寻找的现有函数:

    spark_apply(sdf, function(data) {
       ...
    })
    

    有关详细信息,请参阅sparklyr 文档中的Distributed R

    SparkR

    使用 SparkR 使用 gapply / gapplyCollect

    gapply(df, groupingCols, function(data) {...} schema)
    

    dapply/dapplyCollect

    dapply(df, function(data) {...}, schema)
    

    UDF。参考

    了解详情。

    请注意,与原生 Spark 代码相比,所有解决方案都较差,在需要高性能时应避免使用。

    【讨论】:

    • 函数spark_apply 不合适,因为它在输入中只接受一个参数,并且不保持与mclapply 相同的结构。如果想在 spark data.frame 的行上应用简单的函数,spark_apply 似乎很好。就我而言,我发现很难并行化具有多个输入的函数(例如,将数据帧、向量、字符串、矩阵等作为输入的函数)SparkR 可能更有用,但出于部署目的,由于它不再由 CRAN 维护,因此对我来说不是一个可能的解决方案。
    • 请注意,与原生 Spark 代码相比,所有解决方案都较差,在需要高性能时应避免使用,感谢您的建议。你有一个很好的参考/指南从哪里开始?
    • 请记住 spark_apply 中的上下文参数,您可以通过这种方式传递额外的变量。要传入多个变量(除了简单函数允许的变量),您可以设置 context={c(y
    【解决方案2】:

    sparklyr::spark_apply 现在可以支持将模型等外部变量作为上下文传递。

    这是我在 sparklyr 上运行 xgboost 模型的示例:

    bst <- xgboost::xgb.load("project/models/xgboost.model")
    res3 <- spark_apply(x = ft_union_price %>% sdf_repartition(partitions = 1500, partition_by = "uid"),
                       f = inference_fn,
                       packages = F,
                       memory = F,
                       names = c("uid",
                                   "action_1",
                                   "pred"), 
                       context = {model <- bst})
    

    【讨论】:

      猜你喜欢
      • 2021-12-01
      • 2016-12-25
      • 2019-10-01
      • 2016-02-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-02-25
      • 1970-01-01
      相关资源
      最近更新 更多