【问题标题】:How to pass a third-party column after a GroupBy and aggregation in PySpark DataFrame?如何在 PySpark DataFrame 中的 GroupBy 和聚合之后传递第三方列?
【发布时间】:2021-05-07 18:17:36
【问题描述】:

我有一个 Spark DataFrame,比如df,我需要对其应用 GroupBy col1,通过最大值 col2 聚合并传递 col3 的相应值(与groupBy 或聚合)。最好用一个例子来说明。

df.show()

+-----+-----+-----+
| col1| col2| col3|
+-----+-----+-----+
|    1|  500|  10 |
|    1|  600|  11 |
|    1|  700|  12 |
|    2|  600|  14 |
|    2|  800|  15 |
|    2|  650|  17 |
+-----+-----+-----+

我可以很容易地执行groupBy和聚合以获得col2中每个组的最大值,使用

import pyspark.sql.functions as F

df1 = df.groupBy("col1").agg(
    F.max("col2").alias('Max_col2')).show()

+-----+---------+
| col1| Max_col2|
+-----+---------+
|    1|      700|
|    2|      800|
+-----+---------+

但是,我正在努力并且想做的是另外传递col3的相应值,从而获得下表:

+-----+---------+-----+
| col1| Max_col2| col3|
+-----+---------+-----+
|    1|      700|  12 |
|    2|      800|  15 |
+-----+---------+-----+

有谁知道如何做到这一点?

非常感谢,

马里奥安萨斯

【问题讨论】:

    标签: python apache-spark pyspark apache-spark-sql


    【解决方案1】:

    你可以聚合一个结构的最大值,然后扩展结构:

    import pyspark.sql.functions as F
    
    df2 = df.groupBy('col1').agg(
        F.max(F.struct('col2', 'col3')).alias('col')
    ).select('col1', 'col.*')
    
    df2.show()
    +----+----+----+
    |col1|col2|col3|
    +----+----+----+
    |   1| 700|  12|
    |   2| 800|  15|
    +----+----+----+
    

    【讨论】:

    • 嗨@mck,非常感谢(再次)您的回答。如果我错了,请纠正我:您的代码最初是否比较并找到结构的第一个元素的最大值?如果元素不相等,它将返回具有更高值的结构,否则它将继续比较结构的第二个元素?
    • 好的,我明白了,非常感谢。我会尽快尝试并验证您的答案。谢谢:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-23
    • 2020-01-28
    • 2020-04-29
    • 1970-01-01
    • 2022-01-17
    • 2020-06-14
    相关资源
    最近更新 更多