【问题标题】:Apache spark aggregation: aggregate column based on another column valueApache Spark 聚合:基于另一列值聚合列
【发布时间】:2019-12-10 22:38:42
【问题描述】:

我不确定我是否正确地问了这个问题,也许这就是我到目前为止没有找到正确答案的原因。无论如何,如果它会重复,我会删除这个问题。

我有以下数据:

id | last_updated | count
__________________________
1  | 20190101     | 3
1  | 20190201     | 2
1  | 20190301     | 1 

我想按“id”列按此数据分组,从“last_updated”获取最大值,关于“count”列,我想保留“last_updated”具有最大值的行的值。所以在那种情况下结果应该是这样的:

id | last_updated | count
__________________________
1  | 20190301     | 1 

所以我想它会是这样的:

df
  .groupBy("id")
  .agg(max("last_updated"), ... ("count"))

我可以使用任何函数来根据“last_updated”列获取“计数”。

我使用的是 spark 2.4.0。

感谢您的帮助

【问题讨论】:

    标签: scala apache-spark aggregate


    【解决方案1】:

    你有两个选择,就我的理解而言,第一个更好

    选项 1 对 ID 执行窗口函数,在该窗口函数上创建一个具有最大值的列。然后选择所需列等于最大值的位置,最后删除列并根据需要重命名最大列

    val w  = Window.partitionBy("id")
    
    df.withColumn("max", max("last_updated").over(w))
      .where("max = last_updated")
      .drop("last_updated")
      .withColumnRenamed("max", "last_updated")
    

    选项 2

    您可以在分组后与原始数据框执行连接

    df.groupBy("id")
    .agg(max("last_updated").as("last_updated"))
    .join(df, Seq("id", "last_updated"))
    

    快速示例

    输入

    df.show
    +---+------------+-----+
    | id|last_updated|count|
    +---+------------+-----+
    |  1|    20190101|    3|
    |  1|    20190201|    2|
    |  1|    20190301|    1|
    +---+------------+-----+
    

    输出 选项 1

    import org.apache.spark.sql.expressions.Window
    import org.apache.spark.sql.functions
    
    val w  = Window.partitionBy("id") 
    
    df.withColumn("max", max("last_updated").over(w))
      .where("max = last_updated")
      .drop("last_updated")
      .withColumnRenamed("max", "last_updated")
    
    
    +---+-----+------------+
    | id|count|last_updated|
    +---+-----+------------+
    |  1|    1|    20190301|
    +---+-----+------------+
    

    选项 2

      df.groupBy("id")
          .agg(max("last_updated").as("last_updated")
          .join(df, Seq("id", "last_updated")).show
    
    
        +---+-----------------+----------+
        | id|     last_updated|    count |
        +---+-----------------+----------+
        |  1|         20190301|         1|
        +---+-----------------+----------+
    

    【讨论】:

    • 感谢您的快速回答。完美的。这两种解决方案都对我有用,但是我对选项 1 有点困惑。我以前从未使用过 Window,所以我必须仔细观察它以了解幕后发生的事情。但无论如何谢谢,我会把你的答案标记为正确的。
    • 谢谢。选项 1 它基本上是为指定窗口 (id) 的每个值检索所需列 (las_updated) 的最大值。
    • 我想我现在明白了。谢谢
    猜你喜欢
    • 2018-09-28
    • 2021-09-26
    • 1970-01-01
    • 1970-01-01
    • 2014-06-30
    • 1970-01-01
    • 1970-01-01
    • 2018-12-10
    • 2015-05-17
    相关资源
    最近更新 更多