【问题标题】:Aggregate and Create Array of `Dictionary` object in Spark Dataframe Column在 Spark Dataframe 列中聚合并创建“字典”对象的数组
【发布时间】:2020-05-07 07:53:51
【问题描述】:

我创建了一个玩具火花数据框:

import numpy as np
import pyspark
from pyspark.sql import functions as sf
from pyspark.sql import functions as F

# sc = pyspark.SparkContext()
# sqlc = pyspark.SQLContext(sc)
df = spark.createDataFrame([('csc123','sr1', 'tac1', 'abc'), 
                            ('csc123','sr2', 'tac1', 'abc'), 
                            ('csc234','sr3', 'tac2', 'bvd'),
                            ('csc345','sr5', 'tac2', 'bvd')
                           ], 
                           ['bug_id', 'sr_link', 'TAC_engineer','de_manager'])
df.show()
+------+-------+------------+----------+
|bug_id|sr_link|TAC_engineer|de_manager|
+------+-------+------------+----------+
|csc123|    sr1|        tac1|       abc|
|csc123|    sr2|        tac1|       abc|
|csc234|    sr3|        tac2|       bvd|
|csc345|    sr5|        tac2|       bvd|
+------+-------+------------+----------+

然后我尝试为每个 bug id 聚合并生成 [sr_link, sr_link] 数组

#df = spark.createDataFrame([('row11','row12'), ('row21','row22')], ['colname1', 'colname2'])

df_drop_dup = df.select('bug_id', 'de_manager').dropDuplicates()

df = df.withColumn('joined_column', 
                    sf.concat(sf.col('sr_link'),sf.lit(' '), sf.col('TAC_engineer')))

df_sev_arr = df.groupby("bug_id").agg(F.collect_set("joined_column")).withColumnRenamed("collect_set(joined_column)","sr_array")

df = df_drop_dup.join(df_sev_arr, on=['bug_id'], how='inner')

df.show()

这是输出:

+------+----------+--------------------+
|bug_id|de_manager|            sr_array|
+------+----------+--------------------+
|csc345|       bvd|          [sr5 tac2]|
|csc123|       abc|[sr2 tac1, sr1 tac1]|
|csc234|       bvd|          [sr3 tac2]|
+------+----------+--------------------+

但我真正期望的实际输出是这样的:

+------+----------+----------------------------------------------------------------------+
|bug_id|de_manager|                                                              sr_array|
+------+----------+----------------------------------------------------------------------+
|csc345|       bvd|                                   [{sr_link: sr5, TAC_engineer:tac2}]|
|csc123|       abc|[{sr_link: sr2, TAC_engineer:tac1},{sr_link: sr1, TAC_engineer: tac1}]|
|csc234|       bvd|                                  [{sr_link: sr3, TAC_engineer: tac2}]|
+------+----------+----------------------------------------------------------------------+

因为我希望最终输出可以保存为 JSON 格式,例如:

'bug_id': 'csc123'
'de_manager': 'abc'
'sr_array':
     'sr_link': 'sr2', 'TAC_engineer': 'tac1'
     'sr_link': 'sr1', 'TAC_engineer': 'tac1'

有人可以帮忙吗?抱歉,我对 Spark Dataframe 中的 MapType 非常陌生。

【问题讨论】:

    标签: python dataframe apache-spark pyspark


    【解决方案1】:

    只是根据您的修改了一些功能并添加了新功能 要求。

    第一部分将保持不变。

    from pyspark.sql import functions as F
    
    # sc = pyspark.SparkContext()
    # sqlc = pyspark.SQLContext(sc)
    df = spark.createDataFrame([('csc123','sr1', 'tac1', 'abc'), 
                                ('csc123','sr2', 'tac1', 'abc'), 
                                ('csc234','sr3', 'tac2', 'bvd'),
                                ('csc345','sr5', 'tac2', 'bvd')
                               ], 
                               ['bug_id', 'sr_link', 'TAC_engineer','de_manager'])
    df.show()
    

    我刚刚修改了第二部分。

    >>> df_drop_dup = df.select('bug_id', 'de_manager').dropDuplicates()
    

    修改了WithcolumnRenamed to Alias的重命名函数并添加了to_json and Struct函数以获得所需的输出,并且对数据框命名df--> df1的修改也很少

    >>> df1 = df.withColumn('joined_column', F.to_json(F.struct(F.col('sr_link'), F.col('TAC_engineer'))))
    
    >>> df_sev_arr = df1.groupby("bug_id").agg(F.collect_set("joined_column").alias("sr_array"))
    
    >>> df = df_drop_dup.join(df_sev_arr, on=['bug_id'], how='inner')
    
    >>> df.show(truncate=False)
    +------+----------+----------------------------------------------------------------------------------+
    |bug_id|de_manager|sr_array                                                                          |
    +------+----------+----------------------------------------------------------------------------------+
    |csc345|bvd       |[{"sr_link":"sr5","TAC_engineer":"tac2"}]                                         |
    |csc123|abc       |[{"sr_link":"sr1","TAC_engineer":"tac1"}, {"sr_link":"sr2","TAC_engineer":"tac1"}]|
    |csc234|bvd       |[{"sr_link":"sr3","TAC_engineer":"tac2"}]                                         |
    +------+----------+----------------------------------------------------------------------------------+
    

    如果您有任何与此相关的问题,请告诉我。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-08-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-01-18
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多