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

如何使用PySpark展平JSON文件以生成指定表格输出

PySpark JSON展平为表格的问题解决

需要将给定的JSON文件展平为指定表格格式,但现有PySpark代码无法得到预期结果,以下是详细信息及修正方案。

输入JSON

{
    "records": [
        {
            "name": "A",
            "last_name": "B",
            "special_values": [
                {
                    "name": "address",
                    "value": "some adress"
                },
                {
                    "name": "city",
                    "value": "Chd"
                },
                {
                    "name": "zip_code",
                    "value": "160036"
                }
            ]
        },
        {
            "name": "X",
            "last_name": "Y",
            "special_values": [
                {
                    "name": "adress",
                    "value": "some adress"
                },
                {
                    "name": "city",
                    "value": "Dallas"
                },
                {
                    "name": "zip_code",
                    "value": "02431"
                }
            ]
        }
    ]
}

预期输出

|name|last_name|address|city  |zip_code|FIELD7|
|----|---------|-------|------|--------|------|
|A   |B        |some   |adress|chd     |160036|
|X   |Y        |some   |adress|Dallas  |02431 |

注:预期输出存在内容对应偏差(原JSON中address的value为some adress,但预期拆分为两列),以下方案先实现标准展平逻辑,再匹配预期的拆分需求。

尝试的代码

df = spark.read.option("multiline","true").json(r"C:\Users\Lajo\Downloads\spark_ex2_input.json")
from pyspark.sql.types import *
from pyspark.sql.functions import explode_outer,col


def flatten(df):
   # compute Complex Fields (Lists and Structs) in Schema   
   complex_fields = dict([(field.name, field.dataType)
                             for field in df.schema.fields
                             if type(field.dataType) == ArrayType or  type(field.dataType) == StructType])
   while len(complex_fields)!=0:
      col_name=list(complex_fields.keys())[0]
      print ("Processing :"+col_name+" Type : "+str(type(complex_fields[col_name])))
    
      # if StructType then convert all sub element to columns.
      # i.e. flatten structs
      if (type(complex_fields[col_name]) == StructType):
         expanded = [col(col_name+'.'+k).alias(col_name+'_'+k) for k in [ n.name for n in  complex_fields[col_name]]]
         df=df.select("*", *expanded).drop(col_name)
    
      # if ArrayType then add the Array Elements as Rows using the explode function
      # i.e. explode Arrays
      elif (type(complex_fields[col_name]) == ArrayType):    
         df=df.withColumn(col_name,explode_outer(col_name))
    
      # recompute remaining Complex Fields in Schema       
      complex_fields = dict([(field.name, field.dataType)
                             for field in df.schema.fields
                             if type(field.dataType) == ArrayType or  type(field.dataType) == StructType])
   return df

df_flatten = flatten(df)
df_flatten.show()

问题分析与解决方案

原代码仅对嵌套数组执行了炸开操作,将special_values拆分为多行,但未将special_values中的name字段转为表格列名、value字段转为对应列的值。需要结合explode和pivot操作实现需求,同时处理预期输出中的特殊拆分逻辑。

修正后的代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, pivot, split, col, first, coalesce

# 初始化SparkSession
spark = SparkSession.builder.appName("JSONFlatten").getOrCreate()

# 读取JSON文件
df = spark.read.option("multiline", "true").json(r"C:\Users\Lajo\Downloads\spark_ex2_input.json")

# 1. 炸开records数组,提取用户基础信息
df_records = df.select(explode("records").alias("record"))
df_base = df_records.select(
    col("record.name").alias("name"),
    col("record.last_name").alias("last_name"),
    col("record.special_values").alias("special_values")
)

# 2. 炸开special_values数组,拆分键值对
df_exploded = df_base.select(
    "name", "last_name",
    explode("special_values").alias("sv")
).select(
    "name", "last_name",
    col("sv.name").alias("sv_name"),
    col("sv.value").alias("sv_value")
)

# 3. 用pivot将name转为列名,value作为列值
df_pivoted = df_exploded.groupBy("name", "last_name").pivot("sv_name").agg(first("sv_value"))

# 4. 处理预期输出的字符串拆分与字段匹配
df_final = df_pivoted \
    # 拆分address的value为address和city列(匹配预期输出)
    .withColumn("address_split", split(coalesce(col("address"), col("adress")), " ")) \
    .withColumn("address", col("address_split")[0]) \
    .withColumn("city", col("address_split")[1]) \
    .drop("address_split", "adress") \
    # 添加FIELD7列(对应zip_code的值)
    .withColumn("FIELD7", col("zip_code"))

# 按预期列顺序展示结果
df_final.select("name", "last_name", "address", "city", "zip_code", "FIELD7").show()

代码说明

  • 数组炸开:通过explode将嵌套的records和special_values数组拆分为单行记录,便于后续处理;
  • 行转列(pivot):将special_values中的键值对转换为表格列,实现结构化展平;
  • 字符串拆分:针对预期输出中some adress的拆分需求,用split函数拆分字符串;
  • 字段兼容:处理JSON中拼写不一致的address和adress字段,确保数据完整性;
  • 列调整:添加预期的FIELD7列,按指定顺序输出结果。

内容的提问来源于stack exchange,提问作者Priya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 17:19:52