Spark读取CSV时如何校验分隔符数量,异常行直接抛出异常?
Great question! Spark's built-in CSV reader does default to padding missing/truncated rows with nulls, and there's no out-of-the-box configuration to throw an error directly when the number of delimiters doesn't match the expected count. But you don't have to stick with RDDs—you can implement this validation cleanly using DataFrame APIs, with better integration with Spark's ecosystem.
Option 1: Validate before parsing with text reader + proper CSV parser
The most reliable approach is to first read the file as raw text, validate each line's field count, and then parse valid lines into a DataFrame. This lets you catch invalid rows early, and you can choose to either filter them out or throw an exception immediately.
First, get the expected field count from the header (using a CSV parser to handle quoted fields correctly):
import org.apache.commons.csv.{CSVFormat, CSVParser} import java.io.StringReader // Read the header line to determine expected field count val headerLine = spark.read.text("your/file/path.csv").first().getString(0) val expectedFieldCount = CSVParser.parse(headerLine, CSVFormat.DEFAULT).iterator().next().size()
Then, use a UDF with the same CSV parser to validate each line's field count:
import org.apache.spark.sql.functions._ // UDF to count fields in a CSV line (handles quotes and escaped delimiters) val countCsvFields = udf((line: String) => { val parser = CSVParser.parse(line, CSVFormat.DEFAULT) parser.iterator().next().size() }) // Read raw text, add field count validation val rawTextDF = spark.read.text("your/file/path.csv") .withColumn("field_count", countCsvFields(col("value"))) // Check for invalid rows and throw exception if found val invalidRowCount = rawTextDF.filter(col("field_count") =!= expectedFieldCount).count() if (invalidRowCount > 0) { throw new IllegalArgumentException(s"Error: Found $invalidRowCount rows with incorrect field count (expected $expectedFieldCount fields)") } // Parse valid rows into a structured DataFrame val finalDF = rawTextDF .filter(col("field_count") === expectedFieldCount) .selectExpr("split(value, ',') as fields") .select((0 until expectedFieldCount).map(i => col("fields")(i).alias(s"column_$i")): _*)
Option 2: Post-parsing validation (less reliable)
If you prefer to use Spark's native CSV reader first, you can check for nulls that were added due to truncation—but this is risky, because legitimate null values in your data will be indistinguishable from padding nulls. Still, if your data has no intentional nulls, this could work:
val df = spark.read .option("header", "true") .csv("your/file/path.csv") // Get list of column names val columns = df.columns // Check if any row has nulls (assuming no legitimate nulls exist) val invalidRows = df.filter(columns.map(col(_).isNull).reduce(_ || _)).count() if (invalidRows > 0) { throw new IllegalArgumentException(s"Error: Found $invalidRows rows with missing/truncated fields") }
Key Notes
- Never use simple
split(",")to count fields—this will misinterpret delimiters inside quoted fields (e.g.,"foo,bar",bazshould count as 2 fields, not 3). Always use a proper CSV parser. - The first option is the most robust because it catches invalid rows before parsing, avoiding unnecessary processing of bad data.
- If you're using Python instead of Scala, the logic is identical—you'd use the standard library's
csv.readerinstead of Apache Commons CSV.
内容的提问来源于stack exchange,提问作者shakachuk

