使用Apache Spark SQL将JSON转Parquet时出现_corrupt_record报错
Spark SQL读取JSON转Parquet _corrupt_record报错排查
核心报错原因
当前抛出AnalysisException的直接原因是Spark读取JSON文件时全量解析失败,所有数据都被归入内部默认的损坏记录列_corrupt_record,整个数据集无有效业务字段,触发Spark 2.3版本之后新增的校验规则:禁止对仅包含损坏记录列的原始JSON/CSV数据集直接执行查询操作。
现存问题点
- 代码存在语法错误:Parquet写入行
dr.write().parquet(JSONSAMPLEFILE.parquet ");缺失字符串开头引号,无法正常编译运行,正确写法需传入字符串类型的输出路径。 - Maven依赖版本冲突严重:
引入的Spark核心、Spark SQL依赖为Scala 2.11编译版本(artifactId后缀_2.11),但Spark Cassandra连接器为Scala 2.12编译版本(后缀_2.12),Scala大版本不兼容会直接导致运行时类异常;手动引入的Jackson 2.8.8、Guava 15.0版本远低于Spark 2.4.3、Hadoop 3.3.3内置的依赖版本,会覆盖内置兼容版本引发JSON解析逻辑异常。 - JSON读取配置与文件格式不匹配:
配置了multiline=true时,Spark会将整个文件作为单个完整JSON对象解析,如果目标文件是每行一个JSON对象的JSON Lines格式,该配置会直接导致解析失败;如果JSON文件存在语法错误、带UTF-8 BOM头、首尾存在非JSON特殊字符、传入路径为目录且包含非JSON文件,也会触发全量解析失败。 - Windows路径写法存在兼容风险:使用双反斜杠写相对路径时,容易因IDE工作目录配置错误导致Spark找不到目标文件,将空内容或目录元数据当成JSON解析。
修复方案
- 修正代码逻辑与配置
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import java.io.File; import static org.apache.spark.sql.functions.col; public class Main { public static void main(String[] args) { try { File file = new File("src/main/resources/HadoopResources"); System.setProperty("hadoop.home.dir", file.getAbsolutePath()); SparkSession spark = SparkSession .builder() .appName("Json To Parquet") .config("spark.master", "local") .getOrCreate(); // 用正斜杠写路径避免转义问题,优先替换为实际JSON文件绝对路径排查 String path = "path/JSONSAMPLEFILE.json"; Dataset<Row> dr = spark.read() // 如果是每行一个JSON的格式,把multiline设为false .option("multiline", "true") .option("mode", "PERMISSIVE") .json(path) .cache(); // 打印Schema确认解析结果 dr.printSchema(); // 统计损坏记录数量 long corruptCnt = dr.filter(col("_corrupt_record").isNotNull()).count(); System.out.println("JSON解析损坏记录数:" + corruptCnt); dr.show(); // 无损坏记录时再写入Parquet,修正路径引号问题 if (corruptCnt == 0) { dr.write().parquet("JSONSAMPLEFILE.parquet"); } } catch (Exception exception) { exception.printStackTrace(); } System.out.println("End"); } }
- 清理并统一Maven依赖版本
所有Spark生态组件的Scala编译版本必须完全一致,删除手动引入的Jackson、Netty依赖,这类依赖由Spark内置提供,手动引入版本不兼容会直接导致运行异常;调整Guava版本匹配Hadoop与Spark的版本要求,修正后的依赖配置参考:
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>org.example</groupId> <artifactId>jPFunc</artifactId> <version>1.0-SNAPSHOT</version> <properties> <maven.compiler.source>11</maven.compiler.source> <maven.compiler.target>11</maven.compiler.target> </properties> <dependencies> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-common</artifactId> <version>3.3.3</version> <exclusions> <!-- 排除hadoop自带低版本guava --> <exclusion> <groupId>com.google.guava</groupId> <artifactId>guava</artifactId> </exclusion> </exclusions> </dependency> <!-- 统一使用Scala 2.11版本适配Spark 2.4.3 --> <dependency> <groupId>com.datastax.spark</groupId> <artifactId>spark-cassandra-connector_2.11</artifactId> <version>2.4.3</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.11</artifactId> <version>2.4.3</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.11</artifactId> <version>2.4.3</version> </dependency> <!-- 引入兼容版本guava --> <dependency> <groupId>com.google.guava</groupId> <artifactId>guava</artifactId> <version>27.0-jre</version> </dependency> </dependencies> </project>
- 校验JSON文件合法性
- 若文件为每行一个JSON对象的格式,将读取配置中的
multiline改为false - 用文本编辑器打开JSON文件,检查是否存在引号不匹配、多余逗号等语法错误,删除文件开头的UTF-8 BOM头和首尾无关字符
- 传入的JSON路径必须指向具体文件,不能指向包含多文件的目录,避免读取到非JSON内容
内容的提问来源于stack exchange,提问作者speci5
相关产品推荐
相关产品推荐

