【问题标题】:Aggregate field by ID and retain ID as column in PySpark dataframe按 ID 聚合字段并将 ID 保留为 PySpark 数据框中的列
【发布时间】:2021-05-10 19:16:07
【问题描述】:

基本上,我正在尝试将接受的响应修改为this question,以保留聚合中使用的 id 字段。

如果我的输入数据框如下所示:

id   |  author
a    |  Smith
a    |  Jones
b    |  Kabir
c    |  Chen
c    |  Zhang
c    |  Culver 

理想情况下,我希望我的输出如下所示:

id | authors             | count
a  | Smith, Jones        |   2
b  | Kabir               |   1
c  | Chen, Zhang, Culver |   3

我已经能够使用以下命令非常接近:

    myDF2 = (myDF 
                .groupby("id") 
                .agg(concat_ws(", ", sort_array(collect_list("author"))).alias("authors")) 
                .groupby("authors") 
                .agg(count("id")).alias("count") 
             )

这会产生我想要的聚合,但我不知道如何在输出中保留或添加 id 字段。

【问题讨论】:

  • 只按作者分组,第二组按 id 分组
  • 是否可以使用作者列的size 而不是第二个 groupBy?

标签: apache-spark pyspark


【解决方案1】:

方法 1

在连接列表中的项目之前,您可以使用size对其进行计数


from pyspark.sql import functions as F
myDf2 = myDf.groupBy("id").agg(F.sort_array(F.collect_list('author')).alias('authors'))
myDf2 = myDf2.select(
   F.col('id'),
   F.concat_ws(", ",F.col('authors')).alias('authors'),
   F.size('authors').alias('count')
)

myDf2.show()

输出

+---+-------------------+-----+
| id|            authors|count|
+---+-------------------+-----+
|  c|Chen, Culver, Zhang|    3|
|  b|              Kabir|    1|
|  a|       Jones, Smith|    2|
+---+-------------------+-----+

在我的输出中c 是第一行,如果订单很重要,您可以使用id 订购

myDf2 = myDf2.orderBy(F.col('id'))

方法2

同样,你可以在你的第一个聚合中使用count

from pyspark.sql import functions as F
myDf2 = myDf.groupBy("id").agg(
    F.concat_ws(", ",F.sort_array(F.collect_list('author'))).alias('authors'),
    F.count('author').alias('count'),
)

myDf2.show()

输出

+---+-------------------+-----+
| id|            authors|count|
+---+-------------------+-----+
|  c|Chen, Culver, Zhang|    3|
|  b|              Kabir|    1|
|  a|       Jones, Smith|    2|
+---+-------------------+-----+

设置

myDFData="""
id   |  author
a    |  Smith
a    |  Jones
b    |  Kabir
c    |  Chen
c    |  Zhang
c    |  Culver 
"""

myDf =sparkSession.createDataFrame([{"id":linesplit[0].strip(), "author":linesplit[1].strip()} for linesplit in [line.split("|") for line in myDFData.strip().split("\n")[1:]]])
myDf.show()

输出

+------+---+
|author| id|
+------+---+
| Smith|  a|
| Jones|  a|
| Kabir|  b|
|  Chen|  c|
| Zhang|  c|
|Culver|  c|
+------+---+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-01-29
    • 2021-06-16
    • 2023-03-27
    • 2020-05-06
    • 1970-01-01
    • 1970-01-01
    • 2012-11-04
    • 2019-11-22
    相关资源
    最近更新 更多