【问题标题】:How can I find length of a column in SparkR如何在 SparkR 中找到列的长度
【发布时间】:2016-04-08 00:22:31
【问题描述】:

我正在将纯 R 代码转换为 SparkR 以有效利用 Spark。

我有下面的 CloseDate 列。

CloseDate
2011-01-08
2011-02-07
2012-04-07
2013-04-18
2011-02-07
2010-11-10
2010-12-09
2013-02-18
2010-12-09
2011-03-11
2011-04-10
2013-06-19
2011-04-10
2011-01-06
2011-02-06
2013-04-16
2011-02-06
2015-09-25
2015-09-25
2010-11-10

我想计算日期增加|减少的次数。我有下面的 R 代码可以做到这一点。

dateChange <- function(closeDate, dir){
  close_dt <- as.Date(closeDate)
  num_closedt_out = 0
  num_closedt_in = 0

  for(j in 1:length(close_dt)) 
  {
    curr <- close_dt[j]
    if (j > 1)
      prev <- close_dt[j-1]
    else 
      prev <- curr
    if (curr > prev){
      num_closedt_out = num_closedt_out + 1
    }
    else if (curr < prev){
      num_closedt_in = num_closedt_in + 1
    }
  }
  if (dir=="inc")
    ret <- num_closedt_out
  else if (dir=="dec")
    ret <- num_closedt_in
  ret
} 

我在这里尝试使用 SparkR df$col。由于 spark 懒惰地执行代码,所以在执行过程中我没有得到 length 的值并得到 NaN 错误。

这是我尝试过的修改后的代码。

DateDirChanges <- function(closeDate, dir){
  close_dt <- to_date(closeDate)
  num_closedt_out = 0
  num_closedt_in = 0

  col_len <- SparkR::count(close_dt)
  for(j in 1:col_len) 
  {
    curr <- close_dt[j]
    if (j > 1)
      prev <- close_dt[j-1]
    else 
      prev <- curr
    if (curr > prev){
      num_closedt_out = num_closedt_out + 1
    }
    else if (curr < prev){
      num_closedt_in = num_closedt_in + 1
    }
  }
  if (dir=="inc")
    ret <- num_closedt_out
  else if (dir=="dec")
    ret <- num_closedt_in
  ret
}

在执行此代码期间如何获取列的长度?或者还有其他更好的方法吗?

【问题讨论】:

    标签: r apache-spark sparkr


    【解决方案1】:

    你不能因为Column 根本没有长度。与您在 R 中所期望的不同,列不代表数据,而是 SQL 表达式和特定的数据转换。此外,Spark DataFrame 中的值顺序是任意的,因此您不能简单地环顾四周。

    如果数据可以像您之前的问题中那样进行分区,您可以像我在in the answer to your previous question 中展示的那样使用窗口函数。否则单独使用 SparkR 没有有效的方法来处理这个问题。

    假设有一种方法可以确定顺序(必需)并且您可以对数据进行分区(希望获得合理的性能),那么您所需要的就是这样:

    SELECT
       CAST(LAG(CloseDate, 1) OVER w > CloseDate AS INT) gt,
       CAST(LAG(CloseDate, 1) OVER w < CloseDate AS INT) lt,
       CAST(LAG(CloseDate, 1) OVER w = CloseDate AS INT) eq
    FROM DF
    WINDOW w AS (
      PARTITION BY partition_col ORDER BY order_col
    )
    

    【讨论】:

    • 我想我也可以对数据进行分区。但是我们使用 datediff 来获得所需的输出。但是在这里我需要编写一个自定义函数。就像,它应该检查值是从它的 LAG 增加还是减少,它应该只返回值增加或减少的次数。所以这个自定义函数必须读取数据来创建一个新列。有没有可能这样做?
    • 只要你有办法确定顺序(必需)和分区(为了性能),这很容易。
    • 这与上一个问题完全相同。我只是坚持使用 tempTable 并得到了滞后。在那里我们使用 datediff 来获得差异。这里我们需要写一些类似 getIncrementCount(df$closeDate, df$lagCloseDate) 的东西。在那个函数中,我需要遍历并保持计数。每次当 closeDate 大于 lagCloseDate 时,该计数必须增加一。我确实参考了 SparkR 提供的一些默认函数,但它们都调用了 java 来执行此操作。 R有可能吗?对不起,如果问题太愚蠢。我对 R 和 SparkR 很陌生
    猜你喜欢
    • 2011-06-20
    • 2011-07-22
    • 1970-01-01
    • 2013-09-10
    • 1970-01-01
    • 1970-01-01
    • 2020-04-14
    • 2010-12-06
    • 1970-01-01
    相关资源
    最近更新 更多