Scala作业导入JSON到Hive时如何删除格式无效行 DROPMALFORMED不生效
解决方案
一、DROPMALFORMED不生效的核心原因
- 未开启
multiLine配置:如果你的JSON文件是整文件包裹为一个JSON数组、单个JSON对象跨多行存储,Spark默认按单行解析的模式下无法正常识别结构,DROPMALFORMED逻辑也不会触发,你给出的示例就属于这类数组型JSON结构。 - 配置不兼容:Spark 2.x与3.x的参数配置存在差异,手动指定Schema时未保留坏行记录字段也会导致模式失效。
- 格式错误超出内置校验范围:部分极端语法错误无法被Spark内置的JSON解析器识别为可丢弃的坏行,会直接抛出异常终止作业。
二、可落地的解决步骤
方案1:修正Spark JSON读取配置(优先尝试)
适配数组型/多行JSON的正确配置如下(Scala示例):
import org.apache.spark.sql.SparkSession val spark = SparkSession.builder() .appName("JSON2Hive") .enableHiveSupport() .getOrCreate() // Spark 3.x 版本配置 val df = spark.read .option("multiLine", "true") // 适配跨行JSON/整文件JSON数组结构 .option("mode", "DROPMALFORMED") .option("columnNameOfCorruptRecord", "_corrupt_record") // 显式指定坏行存储字段 .json("你的JSON文件路径") .drop("_corrupt_record") // 校验完成后丢弃坏行列
如果使用Spark 2.x版本,将mode参数替换为.option("dropMalformed", "true")即可。
方案2:自定义坏行过滤(兼容性最高,100%可控)
如果调整配置后仍然报错,直接手动校验JSON合法性再解析,完全可控过滤逻辑:
import com.fasterxml.jackson.databind.ObjectMapper import com.fasterxml.jackson.module.scala.DefaultScalaModule val mapper = new ObjectMapper() mapper.registerModule(DefaultScalaModule) // JSON合法性校验函数 def isValidJson(jsonStr: String): Boolean = { try { mapper.readTree(jsonStr) true } catch { case _: Exception => false } } // 针对整文件JSON数组的处理逻辑 val wholeFileStr = spark.sparkContext.wholeTextFiles("你的JSON文件路径").map(_._2).first() // 移除首尾数组括号,拆分单个JSON对象 val jsonObjects = wholeFileStr.trim.stripPrefix("[").stripSuffix("]").split("\\},\\s*\\n") .map(objStr => if (objStr.endsWith("}")) objStr else s"$objStr}") // 过滤合法JSON后转DataFrame val validJsonRDD = spark.sparkContext.parallelize(jsonObjects).filter(isValidJson) val df = spark.read.json(validJsonRDD)
三、写入Hive
校验过滤完成后直接写入Hive表即可:
df.write.mode("append").saveAsTable("库名.表名")
该逻辑会直接丢弃所有格式非法的JSON条目,不会触发作业中断。
内容的提问来源于stack exchange,提问作者Galay
相关产品推荐
相关产品推荐

