如何用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
相关产品推荐
相关产品推荐

