如何用PySpark生成指定结构的嵌套JSON文件
PySpark 生成指定嵌套JSON结构的解决方案
问题分析
你的代码存在几个核心问题:
- 使用
F.to_json()将结构体转为字符串,导致最终数组中存储的是JSON字符串而非JSON对象,不符合目标格式要求 - 分组后生成的字段名是
values,而非目标结构要求的NewData - 未处理
id字段的类型转换(目标结构中id是字符串类型,原数据为整数)
修正后的完整代码
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, IntegerType, StringType import pyspark.sql.functions as F # 初始化SparkSession spark = SparkSession.builder.appName("NestedJsonOutput").getOrCreate() data = [(1,12,"smith", "uber"), (2,13,"jon","lunch"), (3,15,"jocelyn","rental"), (4,15,"megan","sds") ] schema = StructType([ StructField('id', IntegerType(), True), StructField('age', IntegerType(), True), StructField('number', StringType(), True), StructField('name', StringType(), True) ]) df = spark.createDataFrame(data,schema) # 1. 将id转为字符串类型,匹配目标结构 df = df.withColumn("id", F.col("id").cast(StringType())) # 2. 构造每个元素的结构体,收集为列表后命名为NewData result_df = df.agg(F.collect_list( F.struct("id", "number", "name", "age") ).alias("NewData")) # 3. 输出为符合要求的嵌套JSON文件 # coalesce(1)确保输出单个文件,multiLine=True确保输出为单JSON对象而非每行一个 result_df.coalesce(1).write \ .mode("overwrite") \ .option("multiLine", "true") \ .json("path/to/your/output/directory") # 查看结果内容 result_df.show(truncate=False)
关键说明
- 类型转换:通过
cast(StringType())将id字段从整数转为字符串,完全匹配目标JSON的字段类型 - 结构体收集:直接使用
F.struct()构造每个元素的结构,再用F.collect_list()收集为数组,避免将结构体转为字符串 - 输出配置:
coalesce(1):将数据合并到一个分区,确保输出单个JSON文件(生产环境数据量大时不建议使用,此处为匹配目标单文件结构)multiLine=True:让Spark输出完整的单JSON对象,而非每行一个JSON元素
- 结构匹配:最终
result_df的结构就是{"NewData": [数组元素]},直接输出即可得到目标格式
内容的提问来源于stack exchange,提问作者AutotelicLearner
相关产品推荐
相关产品推荐

