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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 09:10:29