Databricks平台如何将DataFrame转换为指定根层级结构的JSON文件
Databricks 环境下DataFrame转指定嵌套JSON实现方案
基于PySpark实现,可直接在Databricks笔记本运行,将包含OBJECTID、SingleLine列的源DataFrame转换为要求的JSON结构,源字段SingleLine自动映射为目标结构中的Address字段。
完整实现代码
from pyspark.sql import functions as F # 1. 测试源数据构造(如果已经持有生产环境的目标DataFrame,可直接跳过此段逻辑) source_data = [ (1234, "sample Address"), (1, "380 New York St"), (2, "1 World Way") ] source_df = spark.createDataFrame(source_data, schema=["OBJECTID", "SingleLine"]) # 2. 单行结构转换:将每行数据包装为attributes嵌套对象 nested_row_df = source_df.select( F.struct( F.col("OBJECTID"), F.col("SingleLine").alias("Address") ).alias("attributes") ) # 3. 顶层结构组装:将所有行聚合为records数组,完全匹配目标JSON格式 final_json_df = nested_row_df.agg( F.collect_list(F.col("attributes")).alias("records") ) # 4. 结果输出 # 替换为实际存储路径,支持DBFS、挂载的S3/ADLS等合规存储路径 output_target_path = "/dbfs/your/actual/output/path" # 小数据量可加coalesce(1)输出单个JSON文件,大数据量禁止使用,避免Driver节点内存溢出 final_json_df.coalesce(1).write.mode("overwrite").json(output_target_path) # 如需在笔记本内直接预览生成的JSON结果,可执行下行代码 # print(final_json_df.toJSON().collect()[0])
关键说明
- 输出的JSON结构完全匹配规范:顶层为
records数组,数组内每个元素为attributes对象,包含OBJECTID和Address两个字段 - Spark分布式写入默认会生成多个JSON分片,属于正常运行行为,如需合并为单文件仅适合10GB以内的小数据集
- 若使用Scala实现,逻辑完全一致,仅需将API语法调整为对应Scala版本即可
内容的提问来源于stack exchange,提问作者Pysparker
相关产品推荐
相关产品推荐

