如何将PySpark DataFrame的Struct列存储为可被json.load读取的JSON文件
解决Spark DataFrame存储JSON适配Python json.load的方法
要让Spark DataFrame存储的JSON文件能直接用Python的json.load读取,核心是输出标准JSON数组格式(而非Spark默认的JSON Lines每行一个对象格式),同时确保嵌套结构(map、struct、array)转成JSON兼容的格式。以下是具体实现方式:
方法一:小数据量场景(直接在Driver端处理)
如果数据量不大,可把Spark DataFrame的数据拉到Driver端,转换成Python字典列表后用json.dump写入标准JSON:
from pyspark.sql import SparkSession import json # 假设已初始化SparkSession并加载得到DataFrame df spark = SparkSession.builder.appName("JSONExport").getOrCreate() df = spark.read.parquet("your_data_source") # 将Spark Row对象递归转为Python字典/列表结构 data_list = [row.asDict(recursive=True) for row in df.collect()] # 写入标准JSON文件(数组格式) with open("transaction_output.json", "w", encoding="utf-8") as f: json.dump(data_list, f, indent=2)
asDict(recursive=True)会自动处理嵌套结构:
- Spark的
map类型转为JSON对象({"key": value}) struct类型转为嵌套字典array类型转为JSON数组,完全符合JSON规范,json.load可以直接读取整个文件。
方法二:大数据量场景(分布式生成JSON数组)
如果数据量较大,直接collect()会导致Driver内存溢出,可通过Spark RDD操作生成单文件的标准JSON数组:
# 将DataFrame转为每行一个JSON字符串的RDD json_rdd = df.toJSON() # 合并为单个分区并包装成JSON数组格式 final_rdd = json_rdd.coalesce(1).mapPartitions(lambda iter: [f"[{','.join(iter)}]"]) # 写入文件目录(会生成一个part-00000文件) final_rdd.saveAsTextFile("transaction_json_dir")
之后只需将生成的part-00000文件重命名为transaction_output.json即可,json.load可直接读取该文件。
注意事项
- 不要直接使用
df.write.json("output_dir"):这种方式生成的是JSON Lines格式(每行一个独立JSON对象),json.load无法直接读取整个文件,只能逐行用json.loads处理。 - 嵌套结构兼容性:Spark的
map、struct、array类型经过上述转换后,都会自动适配JSON格式,无需额外处理。
内容的提问来源于stack exchange,提问作者datahack
相关产品推荐
相关产品推荐

