Spark中使用groupBy分组后,如何为每个分组生成独立DataFrame?
如何按分组拆分Spark DataFrame为多个子DataFrame?
基于以下Spark DataFrame,我们已按firstname和middlename完成分组,现在需要将每个分组单独拆分为独立的DataFrame:
示例数据准备
data = [('James','','Smith','1991-04-01','M',3000), ('Michael','Rose','','2000-05-19','M',4000), ('Michael','Rose','Rose','1978-09-05','M',4000), ('Jen','Mary','Jones','1967-12-01','F',4000), ('Jen','Mary','Brown','1980-02-17','F',-1) ] columns = ["firstname","middlename","lastname","dob","gender","salary"] df = spark.createDataFrame(data=data, schema = columns)
分组结果
grouped = df.groupBy(['firstname', 'middlename']).count()
分组输出:
+---------+----------+-----+ |firstname|middlename|count| +---------+----------+-----+ | James| | 1| | Michael| Rose| 2| | Jen| Mary| 2| +---------+----------+-----+
解决方案:提取分组键生成子DataFrame
通过提取所有分组的唯一键组合,循环过滤原DataFrame即可生成对应分组的子DataFrame:
# 获取所有分组的键组合并转为列表 group_keys = grouped.select('firstname', 'middlename').collect() # 用字典存储子DataFrame,方便按分组键访问 group_dfs = {} # 遍历每个分组键,生成对应子DataFrame for key in group_keys: fn_val = key.firstname mn_val = key.middlename # 过滤原DataFrame sub_df = df.filter((df.firstname == fn_val) & (df.middlename == mn_val)) # 自定义字典键,避免空字符串影响 df_key = f"{fn_val}_{mn_val}".replace("_", "") if mn_val else fn_val group_dfs[df_key] = sub_df
验证输出
查看James分组的子DataFrame:
group_dfs['James'].show()输出:
+---------+----------+--------+----------+------+------+ |firstname|middlename|lastname| dob|gender|salary| +---------+----------+--------+----------+------+------+ | James| | Smith|1991-04-01| M| 3000| +---------+----------+--------+----------+------+------+查看Michael_Rose分组的子DataFrame:
group_dfs['MichaelRose'].show()输出:
+---------+----------+--------+----------+------+------+ |firstname|middlename|lastname| dob|gender|salary| +---------+----------+--------+----------+------+------+ | Michael| Rose| |2000-05-19| M| 4000| | Michael| Rose| Rose|1978-09-05| M| 4000| +---------+----------+--------+----------+------+------+查看Jen_Mary分组的子DataFrame:
group_dfs['JenMary'].show()输出:
+---------+----------+--------+----------+------+------+ |firstname|middlename|lastname| dob|gender|salary| +---------+----------+--------+----------+------+------+ | Jen| Mary| Jones|1967-12-01| F| 4000| | Jen| Mary| Brown|1980-02-17| F| -1| +---------+----------+--------+----------+------+------+
注意事项
- 该方法适合分组数量较少的场景,若分组数量极大,生成大量子DataFrame会占用过多Driver端资源,需谨慎使用。
- 所有子DataFrame均为原DataFrame的逻辑视图,仅在触发
show()、count()等动作时才会执行实际计算。
内容的提问来源于stack exchange,提问作者Ishan Gote
相关产品推荐
相关产品推荐

