PySpark加载JSON列表时如何避免生成_corrupt_record列
解决Spark JSON解析中
tag为None时的_corrupt_record问题 问题原因
Spark自动推断Schema(inferSchema=true)时,会依据第一条记录的字段类型生成Schema:第一条记录的tag是嵌套对象,被推断为StructType。后续记录中tag为Python的None时,若未正确转为JSON标准的null,或与推断的StructType不兼容,会被判定为格式错误,进入_corrupt_record列。
解决方案
方案1:显式定义Schema(推荐)
通过提前定义Schema,强制指定tag为允许null的嵌套结构,避免自动推断的类型冲突。
from pyspark.sql.types import StructType, StructField, StringType, IntegerType import json # 定义Schema:tag字段为允许null的嵌套结构 schema = StructType([ StructField("name", StringType(), nullable=True), StructField("value", IntegerType(), nullable=True), StructField("tag", StructType([ StructField("property1", StringType(), nullable=True) ]), nullable=True) ]) # 将Python字典列表转为合法JSON字符串(自动把None转为JSON null) json_strings = [json.dumps(item) for item in array_json] # 加载DataFrame df = spark.read\ .option("multiline", True)\ .schema(schema)\ .json(sc.parallelize(json_strings)) df.show(truncate=False)
方案2:预处理JSON数据,确保None转为JSON null
若不想提前定义Schema,需先将Python字典转为标准JSON字符串,让None自动转为JSON的null,保证Spark推断Schema时能识别类型兼容的null值。
import json # 预处理:将Python字典转为合法JSON字符串 json_strings = [json.dumps(item) for item in array_json] # 加载DataFrame df = spark.read\ .option("inferSchema", "true")\ .option("multiline", True)\ .json(sc.parallelize(json_strings)) df.show(truncate=False)
验证结果
两种方案执行后,都能得到期望的DataFrame:
| name | value | tag |
|---|---|---|
| Alex | 2 | {"property1":"value1"} |
| Robert | 2 | null |
内容的提问来源于stack exchange,提问作者Alejandro Alvarez
相关产品推荐
相关产品推荐

