【问题标题】:Sparklyr: how to calculate correlation coefficient between 2 Spark tables?Sparklyr:如何计算 2 个 Spark 表之间的相关系数?
【发布时间】:2017-09-23 18:14:14
【问题描述】:

我有这 2 个 Spark 表:

simx
x0: num 1.00 2.00 3.00 ...
x1: num 2.00 3.00 4.00 ...
...
x788: num 2.00 3.00 4.00 ...

simy
y0: num 1.00 2.00 3.00 ...

在这两个表中,每一列都有相同数量的值。表xy 分别保存到句柄simX_tblsimY_tbl 中。实际数据量很大,可能达到40GB。

我想计算simx 中每一列与simy 的相关系数(比如cor(x0, y0, 'pearson'))。

我到处搜索,我认为没有任何现成的cor 函数,所以我正在考虑使用相关公式本身(就像mentioned in here)。

基于my previous question 中的一个很好的解释,我认为使用mutate_allmutate_each 效率不高,并且为更大的数据大小提供C stack error,所以我考虑使用invoke 来代替调用直接来自Spark 的函数。

到目前为止,我设法到达这里:

exprs <- as.list(paste0("sum(", colnames(simX_tbl),")"))

corr_result <- simX_tbl%>%  
  spark_dataframe() %>% 
  invoke("selectExpr", exprs) %>% 
  invoke("toDF", as.list(colnames(simX_tbl))) %>% 
  sdf_register("corr_result")

计算simx 中每一列的sum。但后来,我意识到我还需要计算simy 表,我不知道如何将这两个表交互在一起(例如,在操作simx 时访问simy)。

有没有什么方法可以更好地计算相关性?或者只是如何与其他 Spark 表交互。

我的 Spark 版本是 1.6.0

编辑: 我尝试使用来自dplyrcombine 函数:

xy_df <- simX_tbl %>% 
  as.data.frame %>%
  combine(as.data.frame(simY_tbl)) %>%
  # convert both table to dataframe, then combine. 
  # It will become list, so need to convert to dataframe again
  as.data.frame 

xydata <- copy_to(sc, xy_df, "xydata") #copy the dataframe into Spark table

但我不确定这是否是一个好的解决方案,因为:

  1. 需要加载到 R 内部的数据框中,我认为这对于大数据不实用
  2. 当尝试head 句柄xydata 时,列名变成所有值的连接

    xydata %>% head
    Source:   query [6 x 790]
    Database: spark connection master=yarn-client app=sparklyr local=FALSE
    

    c_1_67027262134984_2_44919662134984_1_85728542134984_1_49317262134984_
    1 1.670273
    2 2.449197
    3 1.857285
    4 1.493173
    5 1.576857
    6 -5.672155

【问题讨论】:

    标签: r apache-spark dplyr sparkr sparklyr


    【解决方案1】:

    我个人会通过返回the input dataset 来解决它。只是为了记录输入数据已使用 CSV 阅读器加载:

    df <- spark_read_csv(
      sc, path = path, name = "simData", delimiter = " ", 
      header = "false", infer_schema = "false"
    ) %>% rename(y = `_c0`, xs = `_c1`)
    

    看起来或多或少像这样:

          y                                                   xs
      <chr>                                                <chr>
    1 21.66     2.643227,1.2698358,2.6338573,1.8812188,3.8708665
    2 35.15 3.422151,-0.59515584,2.4994135,-0.19701914,4.0771823
    3 15.22  2.8302398,1.9080592,-0.68780196,3.1878228,4.6600842
    

    现在让我们一起处理这两个部分,而不是将数据拆分为多个表:

    exprs <- lapply(
     0:(n - 1), 
     function(i) paste("CAST(xs[", i, "] AS double) AS x", i, sep=""))
    
    df %>% 
      # Convert to native Spark
      spark_dataframe() %>%
      # Split and select xs, but retain y
      invoke("selectExpr", list("y", "split(xs, ',') AS  xs")) %>%
      invoke("selectExpr", c("CAST(y AS DOUBLE)", exprs)) %>%
      # Register table so we can access it from dplyr
      invoke("registerTempTable", "exploded_df")
    

    并申请summarize_each:

    tbl(sc, "exploded_df") %>% summarize_each(funs(corr(., y)), starts_with("x"))
    
    Source:   query [1 x 5]
    Database: spark connection master=local[*] app=sparklyr local=TRUE
    
             x0         x1        x2         x3         x4
          <dbl>      <dbl>     <dbl>      <dbl>      <dbl>
    1 0.8503358 -0.9972426 0.7242708 -0.9975092 -0.5571591
    

    快速健全性检查(yx0yx4 之间的相关性):

    cor(c(21.66, 35.15, 15.22), c(2.643227, 3.422151, 2.8302398))
    
    [1] 0.8503358
    
    cor(c(21.66, 35.15, 15.22), c(3.8708665, 4.0771823, 4.6600842))
    
    [1] -0.5571591
    

    您当然可以先将数据居中:

    exploded <- tbl(sc, "exploded_df")
    
    avgs <- summarize_all(exploded, funs(mean)) %>% as.data.frame()
    center_exprs <- as.list(paste(colnames(exploded ),"-", avgs))
    
    transmute_(exploded, .dots = setNames(center_exprs, colnames(exploded))) %>% 
      summarize_each(funs(corr(., y)), starts_with("x"))
    

    但是it doesn't affect the result:

    Source:   query [1 x 5]
    Database: spark connection master=local[*] app=sparklyr local=TRUE
    
             x0         x1        x2         x3         x4
          <dbl>      <dbl>     <dbl>      <dbl>      <dbl>
    1 0.8503358 -0.9972426 0.7242708 -0.9975092 -0.5571591
    

    如果transmute_summarize_each 都导致了一些问题,我们可以将居中和相关性直接推送到 Spark:

    #Centering
    center_exprs <- as.list(paste(colnames(exploded ),"-", avgs))
    
    exploded %>%  
      spark_dataframe() %>% 
      invoke("selectExpr", center_exprs) %>% 
      invoke("toDF", as.list(colnames(exploded))) %>%
      invoke("registerTempTable", "centered")
    
    centered <- tbl(sc, "centered")
    
    #Correlation
    corr_exprs <- lapply(
      0:(n - 1), 
      function(i) paste("corr(y, x", i, ") AS x", i, sep=""))
    
    centered %>% 
      spark_dataframe() %>% 
      invoke("selectExpr", corr_exprs) %>% 
      invoke("registerTempTable", "corrs")
    
     tbl(sc, "corrs")
    
    Source:   query [1 x 5]
    Database: spark connection master=local[*] app=sparklyr local=TRUE
    
             x0         x1        x2         x3         x4
          <dbl>      <dbl>     <dbl>      <dbl>      <dbl>
    1 0.8503358 -0.9972426 0.7242708 -0.9975092 -0.5571591
    

    中间表当然不是必需的,这可以在我们从数组中提取数据的同时应用。

    【讨论】:

    • 太棒了,陛下!居中部分也会导致C stack error,但是我们可以使用invoke来避免这种情况。为了让您免于小问题,我编辑了您的答案并添加了居中部分,以便其他读者也可以参考。
    猜你喜欢
    • 2015-07-20
    • 1970-01-01
    • 1970-01-01
    • 2020-05-12
    • 1970-01-01
    • 2021-12-31
    • 2021-01-03
    • 1970-01-01
    相关资源
    最近更新 更多