PySpark写入JSON如何保留arrayOfObjects对象数组格式
问题原因
你遇到的格式变化是Spark原生JSON读写机制导致的:
- 你贴出的代码本身存在笔误,正确读JSON的写法是
spark.read.json("/path/source/"),而非df.spark.read() - Spark的
df.write.json()默认输出的是**JSON Lines(行分隔JSON)**格式:每行存储1个独立JSON对象,本身就不支持生成顶层包裹所有记录的数组结构。哪怕你读入的源文件是数组格式,Spark读入时会自动把数组元素展开成独立的行,天然丢失了顶层数组的包裹结构,所以直接读写必然无法保留原格式。
解决方案
根据你是否需要修改文件内容,选择对应方案即可:
方案1:无需修改文件内容,仅做路径拷贝
如果不需要对JSON内的数据做任何转换、过滤等操作,完全没必要走DataFrame读写流程,直接用Hadoop文件系统API做文件拷贝即可,100%保留原文件格式、内容,性能远高于读写DataFrame。
示例代码:
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration() fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) src_path = spark._jvm.org.apache.hadoop.fs.Path("/path/source/") dst_path = spark._jvm.org.apache.hadoop.fs.Path("/path/target/") # 对应overwrite模式:目标路径存在则先删除 if fs.exists(dst_path): fs.delete(dst_path, True) # 执行文件拷贝 spark._jvm.org.apache.hadoop.fs.FileUtil.copy( src_path, fs, dst_path, fs, False, hadoop_conf )
方案2:需要对数据做转换后输出数组格式JSON
如果需要先处理数据再输出,需要手动把所有记录聚合为数组结构再输出,注意该方案会将全量数据聚合到单节点,仅适合数据量不大的场景,数据量过大会引发内存溢出。
示例代码:
from pyspark.sql import functions as F # 正确读入JSON文件 df = spark.read.json("/path/source/") # 此处写你自己的数据处理逻辑,比如过滤、字段选择、关联等 # df = df.filter(F.col("status") == 1).select("id", "name", "create_time") # 将所有行聚合为单个数组,转换为标准JSON字符串 json_output_df = df.agg( F.collect_list(F.struct(*df.columns)).alias("records") ).select(F.to_json("records").alias("value")) # coalesce(1)将数据合并到单个分区,输出单个文本文件 json_output_df.coalesce(1).write.mode("overwrite").text("/path/target/")
输出完成后,可以直接将目标目录下生成的part-*文本文件重命名为.json后缀,就是你需要的顶层数组格式JSON,目录下的_SUCCESS、.crc校验文件可按需删除。
注意事项
- 无数据修改需求时优先选方案1,不会出现字段顺序变化、数值类型自动转换、特殊字符转义等格式偏差问题
- 方案2输出的JSON默认字段顺序按DataFrame的列顺序排列,如果需要严格匹配源文件的字段顺序,读入时可以自定义Schema指定字段顺序,或者在
struct内按顺序指定字段即可
内容的提问来源于stack exchange,提问作者Baobab
相关产品推荐
相关产品推荐

