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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:23:16