You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 21:15:38