PySpark读取CSV指定自定义Schema时FOREIGNKEY列全为null问题
问题根因
Spark CSV数据源在手动传入自定义Schema时,默认按列的物理位置做字段映射,不会自动匹配表头的列名,这是和自动推导Schema场景的核心差异:
- 自动推导Schema时,Spark会先读取表头识别所有列的排列顺序:第1列为
CPRIMARYKEY、第2列为GENDER、第3列为FOREIGNKEY,再逐列扫描值推导类型,因此第3列的FOREIGNKEY能被正常识别为LongType。 - 你自定义的Schema仅定义了2个字段,顺序为第1位
CPRIMARYKEY、第2位FOREIGNKEY,读取时Spark会直接将CSV第1列的值映射给CPRIMARYKEY,将CSV第2列(存储GENDER值的字符串列)映射给FOREIGNKEY字段,字符串值无法转换为LongType就会全部返回null,CSV第3列存储的真实FOREIGNKEY值因为Schema没有定义对应位置的字段,根本不会被加载。
验证&修复方案
快速验证方法
把自定义Schema里的FOREIGNKEY类型临时改为T.StringType()再读取,你会看到该列输出的全部是GENDER字段的内容,即可确认是列位置错位导致的问题。
修复方式
- 方式1:自定义Schema严格对齐CSV的列物理顺序,不需要的列可以定义后在读取完成再筛选,示例代码:
import pyspark.sql.types as T child_schema = T.StructType([ T.StructField("CPRIMARYKEY", T.LongType()), T.StructField("GENDER", T.StringType()), # 按CSV顺序补全中间列占位 T.StructField("FOREIGNKEY", T.LongType()) ]) child_df2 = spark.read.csv("E:\\data\\person.csv", schema=child_schema, multiLine=True, header=True) # 筛选需要的字段即可 child_df2.select("CPRIMARYKEY", "FOREIGNKEY").show()
- 方式2:如果不想手动补全所有中间列,可以先不指定Schema读取数据(或仅开启inferSchema读取基础结构),再通过
select+cast手动将目标列转为LongType,也能达到相同效果。
内容的提问来源于stack exchange,提问作者Deepak_Spark_Beginner
相关产品推荐
相关产品推荐

