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

使用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解析。

修复方案

  1. 修正代码逻辑与配置
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");
    }
}
  1. 清理并统一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>
  1. 校验JSON文件合法性
  • 若文件为每行一个JSON对象的格式,将读取配置中的multiline改为false
  • 用文本编辑器打开JSON文件,检查是否存在引号不匹配、多余逗号等语法错误,删除文件开头的UTF-8 BOM头和首尾无关字符
  • 传入的JSON路径必须指向具体文件,不能指向包含多文件的目录,避免读取到非JSON内容

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.31 11:03:20