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

Spark从JSON文件创建DataFrame的方法及损坏数据报错问题求助

Spark从JSON文件创建DataFrame的方法及损坏数据报错问题求助

嗨,我来帮你搞定这个Spark读取JSON文件时碰到的报错问题~

你遇到的这个pyspark.sql.utils.AnalysisException是Spark 2.3之后的一个限制:当你的查询只引用了内部损坏记录列(默认叫_corrupt_record)时,直接对原始JSON/CSV文件执行这类查询是不被允许的,比如直接过滤或选择这个列的操作。

解决方法其实很简单,按照报错提示的思路来就行:

首先要把解析后的DataFrame缓存或者持久化,之后再进行损坏数据的查询操作,具体步骤如下:

  1. 先读取JSON文件并缓存结果
    如果是直接读取没有自定义Schema的情况:

    df = spark.read.json("你的JSON文件路径").cache()
    

    如果你有预先定义好的Schema(更推荐,解析效率更高),可以这样写:

    from pyspark.sql.types import StructType, StructField, StringType, IntegerType
    
    # 这里根据你的JSON数据结构定义对应Schema
    custom_schema = StructType([
        StructField("id", IntegerType(), nullable=True),
        StructField("username", StringType(), nullable=True),
        # 按需添加其他字段
    ])
    
    df = spark.read.schema(custom_schema).json("你的JSON文件路径").cache()
    
  2. 缓存完成后,就可以正常查询损坏数据了

    • 查看所有损坏的记录内容:
      df.select("_corrupt_record").show(truncate=False)
      
    • 统计损坏记录的数量:
      df.filter(df["_corrupt_record"].isNotNull()).count()
      

为什么要这么做?

Spark 2.3之后做了这个限制,主要是为了避免重复解析原始JSON文件造成的性能损耗。缓存或者持久化解析后的DataFrame后,后续对损坏数据的查询就会基于已经解析好的结果,而不是反复读取解析原始文件。

备注:内容来源于stack exchange,提问作者Va01nita

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.22 12:08:09