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

PySpark读取JSON文件报错:Spark 2.3禁止仅查询_corrupt_record列

PySpark读取JSON报错问题解决

问题核心限制(Spark 2.3+)

自Spark 2.3版本起,禁止直接从原始JSON/CSV文件执行仅引用损坏记录列(默认名为_corrupt_record)的查询,以下写法会触发报错:

错误示例1:

spark.read.schema(schema).json(file).filter($"_corrupt_record".isNotNull).count()

错误示例2:

spark.read.schema(schema).json(file).select("_corrupt_record").show()

正确处理方案

先缓存或持久化解析后的DataFrame,再执行相关查询:

val df = spark.read.schema(schema).json(file).cache()

之后执行查询:

df.filter($"_corrupt_record".isNotNull).count()

测试JSON数据

[
  {
    "student_id": 1,
    "name": "John Doe",
    "age": 18,
    "grade": "A"
  },
  {
    "student_id": 2,
    "name": "Jane Smith",
    "age": 17,
    "grade": "B"
  },
  {
    "student_id": 3,
    "name": "Bob Johnson",
    "age": 19,
    "grade": "C"
  },
  {
    "student_id": 4,
    "name": "Alice Williams",
    "age": 18,
    "grade": "A"
  },
  {
    "student_id": 5,
    "name": "Charlie Brown",
    "age": 17,
    "grade": "B"
  },
  {
    "student_id": 6,
    "name": "Emma Davis",
    "age": 19,
    "grade": "C"
  },
  {
    "student_id": 7,
    "name": "James Miller",
    "age": 18,
    "grade": "A"
  },
  {
    "student_id": 8,
    "name": "Sophie Taylor",
    "age": 17,
    "grade": "B"
  },
  {
    "student_id": 9,
    "name": "David White",
    "age": 19,
    "grade": "C"
  }
]

你的Python代码

mydata = spark.read.json("/original.csv")
mydata.show()

额外注意事项

你用spark.read.json()读取了后缀为.csv的文件,这会让Spark强制按JSON格式解析CSV内容,大概率引发解析错误:

  • 若文件是CSV格式,改用spark.read.csv()读取;
  • 若文件是JSON格式,确保后缀为.json,或显式指定格式:spark.read.format("json").load("/original.json")。

内容的提问来源于stack exchange,提问作者Narenthra Kishore

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 06:50:12