Spark读取CSV时FAILFAST模式未触发异常求助
Spark CSV FAILFAST模式未检测畸形记录的解决方法
核心原因
Spark采用惰性求值机制:你定义的DataFrame读取逻辑只是声明了计算流程,并没有实际执行数据读取和校验。cache()属于转换操作(Transformation),同样不会触发计算,只有遇到**行动操作(Action)**时,才会真正执行数据读取、解析和schema校验。
解决方案
触发实际计算
在定义完DataFrame后,添加一个行动操作(比如count()、show()、collect()等),强制Spark执行读取流程,此时FAILFAST模式才会生效,检测到列数不匹配的畸形记录并抛出异常。修正后的代码示例:
from pyspark.sql.types import StructType, StructField, StringType, DecimalType, DateType sch = StructType( [ StructField("SecurityDescription", StringType(), True), StructField("RIC", StringType(), False), StructField("UniversalClosePrice", DecimalType(18, 9), False), StructField("UniversalClosePriceDate", DateType(), False), StructField("FXIRScalingFactor", StringType(), False), StructField("FileCode", StringType(), False), StructField("BaseCurrencyCode", StringType(), True), ] ) dfR = ( spark.read.format("csv") .option("mode", "FAILFAST") .option("header", "true") .schema(sch) .load(fileName) ) # 触发行动操作,强制执行读取和校验 dfR.count()验证CSV解析配置
确保CSV读取器的基础配置正确,避免误判列数:- 确认分隔符是预期的逗号
option("sep", ","),若之前修改过需改回 - 确保引号处理
option("quote", "\"")处于开启状态,防止将字段内的逗号误判为列分隔符,导致列数计算偏差
- 确认分隔符是预期的逗号
关键说明
- FAILFAST模式的校验逻辑只有在实际解析记录时才会触发,没有行动操作的话,Spark只会停留在执行计划阶段,不会读取文件内容。
cache()仅用于缓存计算结果,本身不会触发计算,必须搭配行动操作才能让缓存和校验逻辑生效。
内容的提问来源于stack exchange,提问作者Wasim Syed
相关产品推荐
相关产品推荐

