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

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:

namevaluetag
Alex2{"property1":"value1"}
Robert2null

内容的提问来源于stack exchange,提问作者Alejandro Alvarez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 18:18:29