使用PySpark内置函数将嵌套JSON解析为多列的实现方案
解决方案:PySpark扁平化嵌套JSON结构并解决类型不匹配错误
错误原因分析
你遇到的Datatype_mismatch.unexpected_input_type错误,是因为from_json()函数要求第一个参数是JSON格式的字符串类型,但你的data字段实际是ARRAY<ARRAY<STRING>>类型——这个字段已经是Spark解析后的数组结构,根本不需要再用from_json解析。
针对示例场景的分步实现
根据你提供的Schema和期望输出,只需逐层展开嵌套数组,提取结构体字段即可:
from pyspark.sql import functions as F # 读取JSON数据(保留你的原有读取代码) df = spark.read.option("multiline", "true").json(json_file_location) # 1. 展开外层Ads数组,得到单个广告结构体 df_explode_ads = df.withColumn("ad_struct", F.explode("Ads")) # 2. 展开Adlist数组,同时提取placement字段 df_explode_adlist = df_explode_ads.withColumn("ad_item", F.explode("ad_struct.Adlist")) \ .withColumn("placement", F.col("ad_struct.placement")) # 3. 提取嵌套字段并整理成目标列 final_df = df_explode_adlist.select( F.col("ID").alias("Id"), "device", "placement", F.col("ad_item.name").alias("name"), F.col("ad_item.subtype").alias("subtype"), F.col("ad_item.category").alias("category") ) final_df.show()
通用扁平化函数(适配大型嵌套Schema)
针对你提到的实际场景(更大的Schema、多层嵌套的Array/Struct),可以用递归函数自动处理所有嵌套结构,无需硬编码字段:
from pyspark.sql import DataFrame from pyspark.sql.types import ArrayType, StructType def flatten_df(df: DataFrame) -> DataFrame: for field in df.schema.fields: # 处理结构体类型:展开所有子字段 if isinstance(field.dataType, StructType): for subfield in field.dataType.fields: df = df.withColumn(f"{field.name}_{subfield.name}", F.col(f"{field.name}.{subfield.name}")) df = df.drop(field.name) # 递归处理新展开的字段(可能仍有嵌套) df = flatten_df(df) # 处理数组类型:先展开数组,再递归处理元素 elif isinstance(field.dataType, ArrayType): df = df.withColumn(f"exploded_{field.name}", F.explode(field.name)) df = df.drop(field.name) df = flatten_df(df) return df # 使用方式:直接传入原始DataFrame flattened_full_df = flatten_df(df) flattened_full_df.show()
这个函数会自动遍历所有字段,逐层展开嵌套的数组和结构体,最终生成完全扁平化的DataFrame。
内容的提问来源于stack exchange,提问作者Codegator
相关产品推荐
相关产品推荐

