Spark Scala读取JSON出现_corrupt_record问题排查求助
问题原因分析与解决方案
你遇到的_corrupt_record问题,核心原因是你的JSON数据格式不符合标准JSON规范,导致Spark的JSON解析器无法正确识别字段。
具体问题点
看你提供的data.json文件,每个JSON对象里的value键没有用双引号包裹:
{"grouping":"group_1", value:5} // 这里的value是非法的键名
根据JSON官方规范,所有对象的键必须是被双引号包裹的字符串。Spark的JSON数据源解析器严格遵循这个规范,所以它无法识别value这个未加引号的键,进而把整行都判定为损坏的记录,最终只会输出_corrupt_record字段。
解决方案
1. 修复JSON数据格式(推荐)
把数据文件里所有的value改成"value",修正后的data.json应该是这样:
{"grouping":"group_1", "value":5} {"grouping":"group_1", "value":6} {"grouping":"group_3", "value":7} {"grouping":"group_2", "value":3} {"grouping":"group_4", "value":2} {"grouping":"group_1", "value":1} {"grouping":"group_2", "value":2} {"grouping":"group_3", "value":3}
修改后重新运行你的代码,Spark就能正确解析出grouping和value两个字段,printSchema()的输出会变成:
root |-- grouping: string (nullable = true) |-- value: long (nullable = true)
2. 额外注意点
另外看你和原书代码的差异:原书里用的是group字段,而你的数据里是grouping,不过这不是导致当前问题的原因,只是后续如果要复用原书的groupBy(expr("myUDF(group)"))逻辑,需要把字段名对应上,或者改成myUDF(grouping)。
验证步骤
- 替换修正后的
data.json文件 - 重新执行
sbt package构建 - 用
spark-submit运行应用,检查printSchema()的输出是否正常
内容的提问来源于stack exchange,提问作者Tim
相关产品推荐
相关产品推荐

