【问题标题】:Pyspark get top two values in column from a group based on orderingPyspark 根据排序从组中获取列中的前两个值
【发布时间】:2021-07-19 21:23:47
【问题描述】:

我正在尝试获取出现在此列表中的前两个计数,按它们出现的最早 log_date。

state   count   log_date
GU      7402    2021-07-19
GU      7402    2021-07-18
GU      7402    2021-07-17
GU      7402    2021-07-16
GU      7397    2021-07-15
GU      7397    2021-07-14
GU      7397    2021-07-13
GU      7402    2021-07-12
GU      7402    2021-07-11
GU      7225    2021-07-10
GU      7225    2021-07-10

在这种情况下,我的预期输出是:

GU      7402    2021-07-16
GU      7397    2021-07-13

这就是我所做的工作,但在一些边缘情况下,计数可能会先下降然后再上升,如上面的示例所示。此代码返回 2021-07-11 作为 count=7402 的最早日期。

    df = df.withColumnRenamed("count", "case_count")
    df2 = df.groupBy(
        "state", "case_count"
    ).agg(
        F.min("log_time").alias("earliest_date")
    )
    df2 = df2.select("state", "case_count", "earliest_date").distinct()

    df = df2.withColumn(
        "last_import_date",
        F.max("earliest_date").over(Window.partitionBy("state"))
    ).withColumn(
        "max_count",
        F.min(
            F.when(
                F.col("earliest_date") == F.col("last_import_date"),
                F.col("case_count")
            )
        ).over(Window.partitionBy("state"))
    )

    df = df.select("state", "max_count", "last_import_date").distinct()

我认为我需要做的是根据 state 和 log_date(desc) 排序选择前两个计数,然后获取每个计数的最小 log_date。我认为 rank() 可能会在这里通过为每个计数获取最高排名来工作,但我对如何将它应用于这种情况感到困惑。无论我尝试什么,我都无法摆脱最后两条 count=7402 记录。也许我忽略了一种更简单的方法?

    df = df.withColumnRenamed(
        "count", "case_count"
    )
   
    df = df.withColumn(
        "rank",
        F.rank().over(
            Window.partitionBy(
                "state", "case_count"
            ).orderBy(
                F.col("state").asc(),
                F.col("log_date").desc()
            )
        )
    ).orderBy(
        F.col("log_date").desc(),
        F.col("state").asc(),
        F.col("rank").desc()
    )

# output
state   count   log_date     rank
GU      7402    2021-07-19   1
GU      7402    2021-07-18   2
GU      7402    2021-07-17   3
GU      7402    2021-07-16   4
GU      7397    2021-07-15   1
GU      7397    2021-07-14   2
GU      7397    2021-07-13   3
GU      7402    2021-07-12   5
GU      7402    2021-07-11   6

【问题讨论】:

    标签: python dataframe apache-spark pyspark


    【解决方案1】:

    你的直觉很正确,这是一个可能的实现

    import pyspark.sql.functions as F
    from pyspark.sql.window import Window
    
    # define some windows for later
    w_date = Window.partitionBy('state').orderBy(F.desc('log_date'))
    w_rn = Window.partitionBy('state').orderBy('rn')
    w_grp = Window.partitionBy('state', 'grp')
    
    df = df\
      .withColumn('rn', F.row_number().over(w_date))\
      .withColumn('changed', (F.col('count') != F.lag('count', 1, 0).over(w_rn)).cast('int'))\
      .withColumn('grp', F.sum('changed').over(w_rn))\
      .filter(F.col('grp') <= 2)\
      .withColumn('min_date', F.col('log_date') == F.min('log_date').over(w_grp))\
      .filter(F.col('min_date') == True)\
      .drop('rn', 'changed', 'grp', 'min_date')
          
    df.show()
    
    +-----+-----+----------+
    |state|count|  log_date|
    +-----+-----+----------+
    |   GU| 7402|2021-07-16|
    |   GU| 7397|2021-07-13|
    +-----+-----+----------+
    

    【讨论】:

    • 我不得不查看滞后,所以在这里肯定学到了一些新东西!谢谢!!!
    • 这是 Stack Overflow 的荣幸,你总能学到新东西 :D 很高兴我能帮上忙!
    猜你喜欢
    • 2021-12-30
    • 1970-01-01
    • 2016-11-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多