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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 20:32:43