【问题标题】:Run length ID in sparklyrsparklyr 中的运行长度 ID
【发布时间】:2017-06-29 17:50:32
【问题描述】:

data.table 提供了一个 rleid 函数,我认为它非常宝贵 - 当观察到的变量发生变化时,它充当一个代码,由其他一些变量排序。

library(dplyr)


tbl = tibble(time = as.integer(c(1, 2, 3, 4, 5, 6, 7, 8)), 
             var  = c("A", "A", "A", "B", "B", "A", "A", "A"))

> tbl
# A tibble: 8 × 2
   time   var
  <int> <chr>
1     1     A
2     2     A
3     3     A
4     4     B
5     5     B
6     6     A
7     7     A
8     8     A

想要的结果是

> tbl %>% mutate(rleid = data.table::rleid(var))
# A tibble: 8 × 3
   time   var rleid
  <int> <chr> <int>
1     1     A     1
2     2     A     1
3     3     A     1
4     4     B     2
5     5     B     2
6     6     A     3
7     7     A     3
8     8     A     3

我想知道是否可以使用sparklyr 提供的工具重现类似的东西。在测试时,我发现我能做的最好的事情就是达到我需要进行填充的点,但后来无法做到。

library(sparklyr)

spark_install(version = "2.0.2")
sc <- spark_connect(master = "local", 
                    spark_home = spark_home_dir())


spk_tbl = copy_to(sc, tbl, overwrite = TRUE)

spk_tbl %>% 
  mutate(var2 = (var != lag(var, 1L, order = time))) %>%  # Thanks @JaimeCaffarel
  mutate(var3 = if(var2) { paste0(time, var) } else { NA })

Source:   query [8 x 4]
Database: spark connection master=local[4] app=sparklyr local=TRUE

   time   var  var2  var3
  <int> <chr> <lgl> <chr>
1     1     A  TRUE    1A
2     2     A FALSE  <NA>
3     3     A FALSE  <NA>
4     4     B  TRUE    4B
5     5     B FALSE  <NA>
6     6     A  TRUE    6A
7     7     A FALSE  <NA>
8     8     A FALSE  <NA>

我尝试过使用SparkR,但我更喜欢sparklyr 接口及其易用性,因此我最好能够在Spark SQL 中执行此操作。

当然,我已经可以通过将数据分成足够小的块来做到这一点,collecting 它,运行一个函数并将其发回。

就上下文而言,我发现rleid 有用的原因是我处理了大量的火车数据,并且能够索引它正在运行的内容非常有用。

感谢您的帮助 阿基尔

【问题讨论】:

  • 我认为你可以使用这个:tbl %&gt;% mutate(rleid = (var != lag(var, 1, default = "asdf"))) %&gt;% mutate(rleid = cumsum(rleid)) 基本上是这个解决方案:stackoverflow.com/a/33510765/2026277
  • @JaimeCaffarel 我没有注意到整洁的cumsum 方法...不幸的是,cumsum 似乎在 Spark-SQL 中不起作用(或者至少我不能这行得通)。 spk_tbl %&gt;% arrange(time) %&gt;% mutate(rleid = (var != lag(var, 1, order = time, default = FALSE))) %&gt;% mutate(rleid = cumsum(rleid))
  • 哦!我错了 - 我只需要先将布尔值转换为 int 。谢谢!考虑添加作为答案+我可以接受。再次感谢

标签: r apache-spark-sql sparklyr


【解决方案1】:

sparklyr 中的可行解决方案如下:

spk_tbl %>% 
  dplyr::arrange(time) %>% 
  dplyr::mutate(rleid = (var != lag(var, 1, order = time, default = FALSE))) %>% 
  dplyr::mutate(rleid = cumsum(as.numeric(rleid)))

【讨论】:

    【解决方案2】:

    试试这个:

    tbl %>% mutate(run = c(0,cumsum(var[-1L] != var[-length(var)])))
    # A tibble: 8 × 3
       time   var   run
      <int> <chr> <dbl>
    1     1     A     0
    2     2     A     0
    3     3     A     0
    4     4     B     1
    5     5     B     1
    6     6     A     2
    7     7     A     2
    8     8     A     2
    

    【讨论】:

    • 谢谢。只是指出在 spark sql 中仍然遇到与原始评论答案相同的铸造问题。我也更喜欢lead/lag soln,因为它强制执行确定性排序!
    猜你喜欢
    • 2021-06-03
    • 1970-01-01
    • 2018-06-09
    • 2022-12-14
    • 2015-05-13
    • 2021-04-24
    • 1970-01-01
    • 2012-08-17
    • 2013-09-27
    相关资源
    最近更新 更多