Spark读取带Schema的CSV时忽略缺失列的实现方法
解决方案
可以实现,核心是利用Spark CSV读取器的列名匹配特性,结合正确的参数设置,无需两次读取文件。具体调整如下:
关键原理
当设置header=True时,Spark会根据CSV文件的表头列名与你定义的Schema字段名进行匹配,而非按位置匹配。对于Schema中存在但源文件里没有的列,Spark会自动填充null,不会将整行标记为坏记录(除非存在其他解析错误,比如类型不匹配)。
修改后的代码
保持你更新后的Schema不变,只需确保读取参数正确(你已大部分设置好):
from pyspark.sql.types import StructType, StructField, StringType, LongType # 包含col_B的目标Schema file_schema = StructType([ StructField("col_A", StringType(), True), StructField("col_B", LongType(), True), # 新增列 StructField("_corrupt_record", StringType(), True) ]) df_src = (spark.read.format('csv') .option("inferSchema", False) .option("header", True) # 关键:按列名匹配 .schema(file_schema) .option("columnNameOfCorruptRecord", "_corrupt_record") .option("mode", "PERMISSIVE") # 默认就是PERMISSIVE,可显式指定 .load(src_file_location) .cache())
效果验证
- 当源文件存在
col_B列时:正常读取该列的值,_corrupt_record为null(无解析错误时)。 - 当源文件不存在
col_B列时:col_B列自动填充null,整行不会被放入_corrupt_record。
额外注意事项
- 确保Schema中新增的
col_B字段设置为可空(True),否则源文件无该列时会抛出异常(你已设置正确)。 - 如果源文件的列名与Schema中
col_B的大小写不一致,可设置spark.sql.caseSensitive = False(默认就是False)来忽略大小写匹配。
内容的提问来源于stack exchange,提问作者xuxu
相关产品推荐
相关产品推荐

