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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 21:48:14