如何让PySpark读取CSV时因类型不匹配抛错而非填充Null
解决PySpark读取CSV时类型不匹配不抛错的问题
方案1:开启ANSI模式(会话级生效)
开启ANSI模式后,PySpark会对非Null值的类型转换失败抛出错误,同时允许合法的Null值(如CSV中的空字段)正常保留。
代码示例
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, IntegerType spark = SparkSession.builder.appName("UnitPriceValidation").getOrCreate() # 开启当前会话的ANSI模式,仅对本次任务生效 spark.conf.set("spark.sql.ansi.enabled", "true") # 定义Schema,UnitPrice指定为IntegerType schema = StructType([ StructField("Id", IntegerType(), nullable=True), StructField("UnitPrice", IntegerType(), nullable=True) ]) # 读取CSV文件,遇到非Null的类型不匹配值会直接抛出异常 df = spark.read.csv("path/to/your/file.csv", schema=schema, header=True) # 执行show等动作时触发校验,错误会立即抛出 df.show()
说明:CSV中类似10.34这类非Null但无法转成整数的值会触发NumberFormatException;而空字段会被正确解析为Null。
方案2:自定义校验(灵活控制错误逻辑)
如果不想全局开启ANSI模式,可先按字符串类型读取字段,手动校验非法值后再转换类型,适合需要自定义错误提示或处理逻辑的场景。
代码示例
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, IntegerType, StringType from pyspark.sql.functions import col spark = SparkSession.builder.appName("UnitPriceCustomCheck").getOrCreate() # 先将UnitPrice按StringType读取,保留原始值 schema = StructType([ StructField("Id", IntegerType(), nullable=True), StructField("UnitPrice", StringType(), nullable=True) ]) df = spark.read.csv("path/to/your/file.csv", schema=schema, header=True) # 筛选出非Null但无法转换为整数的非法记录 invalid_rows = df.filter( col("UnitPrice").isNotNull() & col("UnitPrice").cast(IntegerType()).isNull() ) # 存在非法记录时主动抛出错误,附带具体行信息 if invalid_rows.count() > 0: error_detail = f"发现{invalid_rows.count()}条非法UnitPrice记录:\n{invalid_rows.collect()}" raise ValueError(error_detail) # 将合法值转换为IntegerType,保留合法Null df = df.withColumn("UnitPrice", col("UnitPrice").cast(IntegerType())) df.show()
说明:精准定位非法行并抛出自定义错误,合法Null不受影响,还可根据需求扩展非法行的处理逻辑(如写入错误日志)。
内容的提问来源于stack exchange,提问作者PaulM
相关产品推荐
相关产品推荐

