如何在PySpark中写入多行JSON格式记录?
PySpark生成标准JSON数组格式文件的解决方案
PySpark默认的df.write.json()输出的是JSON Lines格式(每行一个独立JSON对象),要生成带方括号、记录间有逗号且每行字段换行的标准JSON数组,可按以下两种场景处理:
场景1:小数据集(数据量不大,可加载到Driver内存)
直接将数据收集到Driver端,用Python原生json库格式化输出:
import json # 将DataFrame转为Python字典列表 data_list = df.toJSON().map(lambda json_str: json.loads(json_str)).collect() # 生成带缩进的JSON数组字符串 formatted_json = json.dumps(data_list, indent=2) # 写入目标文件 with open("output.json", "w", encoding="utf-8") as f: f.write(formatted_json)
场景2:大数据集(避免Driver内存溢出)
先将每条记录格式化为带缩进的JSON字符串,再拼接成数组格式,减少Driver端的内存占用:
import json from pyspark.sql.functions import udf, struct from pyspark.sql.types import StringType # 定义UDF:将每行数据转为带缩进的JSON字符串 def format_single_row(row): return json.dumps(row.asDict(), indent=2) format_row_udf = udf(format_single_row, StringType()) # 生成每行格式化后的JSON字符串 formatted_df = df.select(format_row_udf(struct(df.columns)).alias("formatted_str")) # 拼接所有行,生成JSON数组 json_array_content = "[" + ",\n".join(formatted_df.select("formatted_str").rdd.flatMap(lambda x: x).collect()) + "]" # 写入文件 with open("output.json", "w", encoding="utf-8") as f: f.write(json_array_content)
补充说明
- 若数据量极大(远超Driver内存上限),建议先输出JSON Lines格式,再用外部工具(如
jq)转换为JSON数组:jq -s . input.json > output.json indent=2参数控制字段缩进的空格数,可根据需求调整。
内容的提问来源于stack exchange,提问作者aQ123
相关产品推荐
相关产品推荐

