【问题标题】:How to split data into groups in pyspark如何在pyspark中将数据分组
【发布时间】:2020-08-01 11:03:57
【问题描述】:

我需要在时间序列数据中找到组。

数据样本

我需要根据valueday 输出列group

我尝试过使用滞后、领先和行号,但结果一无所获。

【问题讨论】:

    标签: sql select pyspark window-functions gaps-and-islands


    【解决方案1】:

    PySpark 方法。使用 lag 查找组的端点,在此 lag 上执行 incremental sum 以获得 groups strong>、add 1 到群组获取您的desired groups.

    from pypsark.sql.window import Window
    from pyspark.sql import functions as F
    
    w1=Window().orderBy("day")
    df.withColumn("lag", F.when(F.lag("value").over(w1)!=F.col("value"), F.lit(1)).otherwise(F.lit(0)))\
      .withColumn("group", F.sum("lag").over(w1) + 1).drop("lag").show()
    
    #+-----+---+-----+
    #|value|day|group|
    #+-----+---+-----+
    #|    1|  1|    1|
    #|    1|  2|    1|
    #|    1|  3|    1|
    #|    1|  4|    1|
    #|    1|  5|    1|
    #|    2|  6|    2|
    #|    2|  7|    2|
    #|    1|  8|    3|
    #|    1|  9|    3|
    #|    1| 10|    3|
    #|    1| 11|    3|
    #|    1| 12|    3|
    #|    1| 13|    3|
    #+-----+---+-----+
    

    【讨论】:

      【解决方案2】:

      您似乎想在每次值更改时增加组。如果是这样,这就是一种孤岛问题。

      这是一种使用lag() 和累积sum() 的方法:

      select
          value,
          day,
          sum(case when value = lag_value then 0 else 1 end) over(order by day) grp
      from (
          select t.*, lag(value) over(order by day) lag_value
          from mytable t
      ) t
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2022-12-17
        • 2021-12-03
        • 2020-06-10
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-07-10
        • 2019-12-07
        相关资源
        最近更新 更多