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

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)

解决方案

核心问题分析

  1. 顶层结构不匹配:自定义Schema直接定义了目标字段,但原始JSON的目标数据嵌套在_embedded.items数组中,顶层只有_embedded和_links两个字段,导致Spark无法匹配到数据,返回全NULL。
  2. 字段类型错误:多个字段被错误定义为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 19:50:33