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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 05:43:16