Spark读取指定Schema的CSV时如何丢弃格式错误的行?
解决Spark DataSet加载CSV时过滤不符合Schema的行
这个问题我之前也碰到过,其实Spark本身就提供了几种便捷的方式来处理这类不符合Schema的行,下面给你两种常用的解决方案:
方案一:加载后过滤null值
当你显式指定Schema读取CSV时,Spark默认的mode是PERMISSIVE——这种模式下,不符合类型定义的字段会被转为null。所以我们只需要过滤掉目标列是null的行即可:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.{StructType, StructField, DataTypes} val schema = StructType(StructField("col", DataTypes.DoubleType) :: Nil) // 先按原方式加载数据,不符合的列会变成null val ds = spark.read.format("csv") .option("delimiter", "\t") .schema(schema) .load("f.csv") // 过滤掉col为null的行 val filteredDs = ds.filter($"col".isNotNull)
这种方式的好处是,你可以先保留所有原始数据(包括有问题的行),后续如果需要分析错误数据也方便处理。
方案二:读取时直接丢弃不符合的行
如果你不需要保留错误数据,更高效的方式是在读取阶段就指定DROPMALFORMED模式,Spark会自动丢弃整个不符合Schema的行:
val schema = StructType(StructField("col", DataTypes.DoubleType) :: Nil) val filteredDs = spark.read.format("csv") .option("delimiter", "\t") .schema(schema) .option("mode", "DROPMALFORMED") // 关键配置 .load("f.csv")
这里顺便提一下Spark CSV读取的三种模式,方便你根据场景选择:
PERMISSIVE(默认):将不符合的数据转为null,保留所有行DROPMALFORMED:丢弃所有包含不符合Schema字段的行FAILFAST:一旦遇到不符合的数据,直接抛出异常终止读取
内容的提问来源于stack exchange,提问作者Zhe Hou
相关产品推荐
相关产品推荐

