Scala中将DataFrame转换为指定嵌套JSON的实现求助
将扁平DataFrame转换为指定嵌套JSON格式
针对你的需求,直接使用toJson无法生成目标嵌套结构——原DataFrame是扁平多行结构,同一item的分散数据需要合并,还需按层级构造嵌套结构体。以下是基于PySpark的实现步骤:
1. 导入依赖
from pyspark.sql import functions as F
2. 分组聚合,收集嵌套数组
按顶层唯一标识(id*1、name*1等)分组,把同一item对应的topping、batter数据收集为结构体数组:
grouped_df = df.groupBy("id*1", "name*1", "ppu*1", "type*1") \ .agg( # 收集topping数组,每个元素是含id/type的结构体 F.collect_list(F.struct( F.col("toppings*1->id*2").alias("id"), F.col("toppings*1->type*2").alias("type") )).alias("topping"), # 收集batter数组,每个元素是含id/type的结构体 F.collect_list(F.struct( F.col("batters*1->batter*2->id*3").alias("id"), F.col("batter*2->type*3").alias("type") )).alias("batter") )
3. 构造层级嵌套结构
把batter数组包装到batters对象中,并修正顶层列名:
nested_df = grouped_df \ # 构建batters层级:{"batters": {"batter": [...]}} .withColumn("batters", F.struct(F.col("batter").alias("batter"))) \ .drop("batter") \ # 去掉列名后缀*1,匹配目标JSON的字段名 .withColumnRenamed("id*1", "id") \ .withColumnRenamed("name*1", "name") \ .withColumnRenamed("ppu*1", "ppu") \ .withColumnRenamed("type*1", "type")
4. 构造最外层JSON结构
将所有item数据包装到指定的外层结构中:
# 先把item转为数组,再嵌套到items对象里 final_df = nested_df.select(F.array(F.struct(*nested_df.columns)).alias("item")) \ .select(F.struct(F.col("item")).alias("items"))
5. 转换为目标JSON
最后将DataFrame转为符合要求的JSON字符串:
target_json = final_df.toJSON().collect()[0]
执行后生成的JSON会完全匹配你给出的嵌套格式,若存在多条不同id的item,item数组会自动包含所有条目。
内容的提问来源于stack exchange,提问作者Samatha
相关产品推荐
相关产品推荐

