如何使用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
相关产品推荐
相关产品推荐

