如何将DataFrame按id分组并转换为指定嵌套JSON格式?
解决Spark DataFrame嵌套聚合问题
原始DataFrame
| id | type | application | number |
|---|---|---|---|
| 1 | test | spark | 2 |
| 1 | test | kafka | 1 |
| 1 | test1 | spark | 2 |
| 2 | test2 | kafka | 1 |
| 2 | test2 | kafka | 1 |
| 3 | test | spark | 2 |
期望输出
| id | type |
|---|---|
| 1 | {"test":{"spark":2,"kafka":1},"test1":{"spark":2}} |
| 2 | {"test2":{"kafka":1}} |
| 3 | {"test":{"spark":2}} |
实现代码(PySpark)
from pyspark.sql import functions as F # 假设原始DataFrame名为df # 1. 去重/聚合相同(id, type, application)组合的记录,避免重复计算 agg_step1 = df.groupBy("id", "type", "application") \ .agg(F.first("number").alias("number")) # 若需求和可替换为F.sum,根据实际需求调整 # 2. 按(id, type)聚合,生成{application: number}格式的内层Map agg_step2 = agg_step1.groupBy("id", "type") \ .agg(F.map_from_entries(F.collect_list(F.struct("application", "number"))).alias("app_map")) # 3. 按id聚合,生成{type: {application: number}}格式的嵌套Map final_df = agg_step2.groupBy("id") \ .agg(F.map_from_entries(F.collect_list(F.struct("type", "app_map"))).alias("type")) # 查看结果 final_df.show(truncate=False)
代码说明
- 第一步聚合:处理重复的
(id, type, application)组合(比如id=2的两条重复数据),确保每个组合仅保留一条记录。 - 生成内层Map:用
struct将application和number打包为结构体,通过collect_list收集同组数据,再用map_from_entries转换为键值对Map。 - 生成外层嵌套Map:复用上述逻辑,将
type与对应的内层Map打包后聚合,最终得到期望的嵌套Map格式。
内容的提问来源于stack exchange,提问作者Sachin
相关产品推荐
相关产品推荐

