PySpark自定义Schema读取JSON全为NULL的Spark原生解决方案求助
问题:自定义Schema读取JSON时Spark DataFrame返回全NULL值
使用自定义Schema读取JSON文件时,Spark DataFrame所有字段值均为NULL。已知问题源于实际数据类型与自定义Schema不匹配,要求采用Spark原生方式解决,禁止使用Python JSON模块的with open方法处理文件。
用户提供的代码
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, BooleanType, ArrayType, MapType, IntegerType spark = SparkSession \ .builder \ .appName("JSON test") \ .getOrCreate() schema = StructType([StructField("_links", MapType(StringType(), MapType(StringType(), StringType()))), StructField("identifier", StringType()), StructField("enabled", BooleanType()), StructField("family", StringType()), StructField("categories", ArrayType(StringType())), StructField("groups", ArrayType(StringType())), StructField("parent", StringType()), StructField("values", MapType(StringType(), ArrayType(MapType(StringType(), StringType())))), StructField("created", StringType()), StructField("updated", StringType()), StructField("associations", MapType(StringType(), MapType(StringType(), ArrayType(StringType())))), StructField("quantified_associations", MapType(StringType(), IntegerType())), StructField("metadata", MapType(StringType(), StringType()))]) df = spark.read.format("json") \ .schema(schema) \ .load(f'/mnt/bronze/products/**/*.json') df.display()
原始JSON结构
root |-- _embedded: struct (nullable = true) | |-- items: array (nullable = true) | | |-- element: struct (containsNull = true) | | | |-- _links: struct (nullable = true) | | | | |-- self: struct (nullable = true) | | | | | |-- href: string (nullable = true) | | | |-- associations: struct (nullable = true) | | | | |-- ERP_PIM: struct (nullable = true) | | | | | |-- groups: array (nullable = true) | | | | | | |-- element: string (containsNull = true) | | | | | |-- product_models: array (nullable = true) | | | |-- categories: array (nullable = true) | | | | |-- element: string (containsNull = true) | | | |-- created: string (nullable = true) | | | |-- enabled: boolean (nullable = true) | | | |-- family: string (nullable = true) | | | |-- groups: array (nullable = true) | | | | |-- element: string (containsNull = true) | | | |-- identifier: string (nullable = true) | | | |-- metadata: struct (nullable = true) | | | | |-- workflow_status: string (nullable = true) | | | |-- parent: string (nullable = true) | | | |-- updated: string (nullable = true) | | | |-- values: struct (nullable = true) | | | | |-- Contrex_table: array (nullable = true) | | | | | |-- element: struct (containsNull = true) | | | | | | |-- data: string (nullable = true) | | | | | | |-- locale: string (nullable = true) | | | | | | |-- scope: string (nullable = true) | | | | |-- UFI_Table: array (nullable = true) | | | | | |-- element: struct (containsNull = true) |-- _links: struct (nullable = true) | |-- first: struct (nullable = true) | | |-- href: string (nullable = true) | |-- next: struct (nullable = true) | | |-- href: string (nullable = true) | |-- self: struct (nullable = true) | | |-- href: string (nullable = true)
解决方案
核心问题分析
- 顶层结构不匹配:自定义Schema直接定义了目标字段,但原始JSON的目标数据嵌套在
_embedded.items数组中,顶层只有_embedded和_links两个字段,导致Spark无法匹配到数据,返回全NULL。 - 字段类型错误:多个字段被错误定义为
MapType,实际是StructType(比如_links、associations、metadata、values)。
修正后的代码
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, BooleanType, ArrayType, IntegerType from pyspark.sql.functions import explode spark = SparkSession \ .builder \ .appName("JSON test") \ .getOrCreate() # 定义匹配原始JSON的完整Schema item_schema = StructType([ StructField("_links", StructType([ StructField("self", StructField("href", StringType())) ])), StructField("associations", StructType([ StructField("ERP_PIM", StructType([ StructField("groups", ArrayType(StringType())), StructField("product_models", ArrayType(StringType())) ])) ])), StructField("categories", ArrayType(StringType())), StructField("created", StringType()), StructField("enabled", BooleanType()), StructField("family", StringType()), StructField("groups", ArrayType(StringType())), StructField("identifier", StringType()), StructField("metadata", StructType([ StructField("workflow_status", StringType()) ])), StructField("parent", StringType()), StructField("updated", StringType()), StructField("values", StructType([ StructField("Contrex_table", ArrayType(StructType([ StructField("data", StringType()), StructField("locale", StringType()), StructField("scope", StringType()) ]))), StructField("UFI_Table", ArrayType(StructType([ # 根据实际JSON结构补充字段,此处为示例占位 ]))) ])) ]) top_level_schema = StructType([ StructField("_embedded", StructType([ StructField("items", ArrayType(item_schema)) ])), StructField("_links", StructType([ StructField("first", StructField("href", StringType())), StructField("next", StructField("href", StringType())), StructField("self", StructField("href", StringType())) ])) ]) # 读取完整结构的JSON df_full = spark.read.format("json") \ .schema(top_level_schema) \ .load('/mnt/bronze/products/**/*.json') # 展开items数组,提取目标数据 df = df_full.select(explode("_embedded.items").alias("item")) \ .select("item.*") df.display()
关键修正说明
- 顶层Schema匹配:先定义包含
_embedded和_links的顶层Schema,确保Spark能正确读取整个JSON结构。 - 展开嵌套数组:用
explode函数将_embedded.items数组展开,把每个数组元素作为单独的行。 - 修正字段类型:将所有实际为嵌套结构体的字段从
MapType改为StructType,严格匹配原始JSON的层级结构。 - 保留原生读取方式:全程使用Spark的
read.json接口,未使用Python JSON模块处理文件,符合要求。
内容的提问来源于stack exchange,提问作者Greencolor
相关产品推荐
相关产品推荐

