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
相关产品推荐
相关产品推荐

