Spark读取JSON文件出现_corrupt_record列问题求助
解决Spark读取JSON时出现
_corrupt_record的问题 我来帮你搞定这个问题——从你给出的错误输出和代码来看,核心问题是你的JSON数据是外层包裹数组的格式,但当前的读取配置没完全匹配,导致Spark把数组的开头符号[当成了无效的独立记录,从而生成了_corrupt_record列。
问题原因拆解
你设置了option("multiline", "true"),这个参数的作用是允许单个JSON对象跨多行存储,但你的JSON文件是一个完整的数组(外层用[]包裹多个对象)。当Spark按这种模式读取时,会把文件中的每一行当作一个独立的JSON对象来解析:
- 第一行的
[显然不是有效的JSON对象,直接被标记为损坏记录 - 后续行的JSON对象因为缺少外层结构,被部分解析,导致字段值不全或为
null
针对性解决方案
根据你的JSON文件格式,分两种场景处理:
场景1:JSON数组紧凑存储(或每个对象单行)
如果你的JSON文件内容是紧凑的一行,或者每个对象单独占一行(外层仍有数组括号),直接去掉multiline参数即可——Spark默认支持解析JSON数组,会自动把数组中的每个对象当作一条数据记录:
val mdf = sparkSession.read.json("multi2.json") mdf.show(false)
场景2:JSON对象跨多行存储(必须用multiline)
如果你的单个JSON对象本身跨了多行(比如对象内的字符串、嵌套结构有换行),需要先把整个文件读取为单个字符串,再通过Schema解析数组:
- 先定义匹配你数据的Schema:
import org.apache.spark.sql.types._ import org.apache.spark.sql.functions._ val dataSchema = StructType(Seq( StructField("string", StringType), StructField("int", IntegerType), StructField("array", ArrayType(IntegerType)), StructField("dict", StructType(Seq(StructField("key", StringType)))) )) // 因为外层是数组,所以要包裹一层ArrayType val arraySchema = ArrayType(dataSchema)
- 读取整个文件为文本,再解析数组并展开成单条记录:
// 用wholetext=true读取整个文件为单行字符串 val textDf = sparkSession.read.option("wholetext", true).text("multi2.json") // 解析JSON数组,再展开成单条记录 val parsedDf = textDf .select(from_json(col("value"), arraySchema).alias("data_list")) .select(explode(col("data_list")).alias("record")) .select("record.*") parsedDf.show(false)
验证建议
你可以先查看一下multi2.json的实际格式:
cat multi2.json
如果输出是紧凑的一行数组,用场景1的方案即可;如果是多行排版的数组(每个对象跨多行),则用场景2的方案。
内容的提问来源于stack exchange,提问作者More Than Five
相关产品推荐
相关产品推荐

