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

如何用Python(Synapse Spark)将CSV转换为指定层级的JSON结构?

CSV转指定层级JSON(Synapse Spark环境)

我正在处理SQL生成的CSV文件,文件包含分析所需的所有字段与计算结果,目标是将其转换为符合指定层级结构的JSON文件,适配自定义Function App的格式要求。当前使用Synapse Spark Notebook进行格式化,Data1与Data2通过外键关联但需作为独立数组输出,目前Data1的格式已正确,Data2的格式化存在问题,无法提供源数据集。

已编写的Python代码:

df = spark.read.option('header', 'true') 
               .option('delimiter', ',') 
               .csv(read_path)

df2 = (
    df.groupBy("Data1Id")
    .agg(collect_set('PreviousData1').alias('PreviousData1'))
)

所需的JSON格式:

{
   "Data1":[
      {
         "Data1Id":"id1",
         "PreviousData1":[
            {
               "Id":"PreviousId",
               "PreviousShifts":[
               ]
            }
         ],
         "FutureData1":[
            {
               "Id":"futureId",
               "FutureData1":[
                  
               ],
               "PreviousData1":[
               ]
            }
         ]
      }
   ],
 "Data2":[
      {
         "Data2Id":"id2",
         "Function2":[
            {
               "FunctionName":"function",
               "Value":"3"
            }
         ]
      }
   ]
}

解决方案

1. 修复Data2的嵌套结构

Data2需要生成Function2嵌套数组,需先将FunctionName和Value打包为结构体,再按Data2Id分组收集为数组:

from pyspark.sql import functions as F

# 处理Data2的层级结构
df_data2 = df.groupBy("Data2Id")\
    .agg(
        F.collect_set(
            F.struct(
                F.col("FunctionName").alias("FunctionName"),
                F.col("Value").alias("Value")
            )
        ).alias("Function2")
    )

2. 完善Data1的完整嵌套(若需补充)

如果Data1的FutureData1结构还未完善,用同样的结构体打包+分组聚合方式处理:

# 处理Data1的完整嵌套结构
df_data1 = df.groupBy("Data1Id")\
    .agg(
        # 生成PreviousData1嵌套数组
        F.collect_set(
            F.struct(
                F.col("PreviousId").alias("Id"),
                F.collect_set(F.col("PreviousShift")).alias("PreviousShifts")
            )
        ).alias("PreviousData1"),
        # 生成FutureData1嵌套数组
        F.collect_set(
            F.struct(
                F.col("futureId").alias("Id"),
                F.collect_set(F.col("FutureField")).alias("FutureData1"),
                F.collect_set(F.col("PrevField")).alias("PreviousData1")
            )
        ).alias("FutureData1")
    )

注意:代码中的字段名(如PreviousShift、FutureField)需替换为CSV中实际对应的列名。

3. 合并生成最终JSON

将处理好的Data1和Data2数据集转换为JSON格式,组合成顶层结构:

# 将DataFrame转为JSON字符串列表
data1_json_list = df_data1.toJSON().collect()
data2_json_list = df_data2.toJSON().collect()

# 组合成目标JSON结构
final_json = {
    "Data1": [eval(json_str) for json_str in data1_json_list],
    "Data2": [eval(json_str) for json_str in data2_json_list]
}

# 格式化输出或保存
import json
print(json.dumps(final_json, indent=2))

# 保存到Synapse可访问的路径(示例为ADLS路径)
with open("/dbfs/mnt/output/final_result.json", "w") as f:
    json.dump(final_json, f, indent=2)

关键注意事项

  • 所有结构体中的字段名必须与CSV实际列名一一对应,需根据源数据调整。
  • 若需保留重复记录,将collect_set替换为collect_list。
  • Synapse中保存文件时,需使用正确的ADLS路径格式(如/dbfs/mnt/...)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 21:53:25