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

如何将DataFrame按id分组并转换为指定嵌套JSON格式?

解决Spark DataFrame嵌套聚合问题

原始DataFrame

idtypeapplicationnumber
1testspark2
1testkafka1
1test1spark2
2test2kafka1
2test2kafka1
3testspark2

期望输出

idtype
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)

代码说明

  1. 第一步聚合:处理重复的(id, type, application)组合(比如id=2的两条重复数据),确保每个组合仅保留一条记录。
  2. 生成内层Map:用struct将application和number打包为结构体,通过collect_list收集同组数据,再用map_from_entries转换为键值对Map。
  3. 生成外层嵌套Map:复用上述逻辑,将type与对应的内层Map打包后聚合,最终得到期望的嵌套Map格式。

内容的提问来源于stack exchange,提问作者Sachin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 09:13:20