PySpark读取CSV自定义Schema无法设置nullable=false问题
问题:PySpark读取CSV时自定义Schema设置nullable=False不生效
我在加载CSV文件时,自定义Schema将id列的nullable设为false,但执行后打印Schema发现该列仍显示nullable=true。检查数据后确认id列没有空值,预期nullable应该是false。
代码
from pyspark.sql.types import * udemy_comments_file = '/Users/harbeerkadian/Documents/workspace/learn-spark/source_data/udemy/comments_spark.csv' schema = StructType([StructField("id",StringType(),False), StructField("course_id",StringType(),True), StructField("rate",DoubleType(),True), StructField("date",TimestampType(),True), StructField("display_name",StringType(),True), StructField("comment",StringType(),True), StructField("new_id",StringType(),True)]) comments_df = spark.read.format('csv').option('header', 'true').schema(schema).load(udemy_comments_file) comments_df.printSchema() print("non null record count for id", comments_df.filter(comments_df.id.isNull()).count())
输出
root |-- id: string (nullable = true) |-- course_id: string (nullable = true) |-- rate: double (nullable = true) |-- date: timestamp (nullable = true) |-- display_name: string (nullable = true) |-- comment: string (nullable = true) |-- new_id: string (nullable = true) non null record count for id 0
原因分析
Spark的CSV数据源读取器会忽略自定义Schema中的nullable=false设置。这是因为CSV是无结构的文本格式,Spark无法在读取阶段保证某列绝对不存在空值(即便当前数据无空值,后续新增数据仍可能出现),为了容错,会强制将所有CSV加载的列标记为nullable=true,不受自定义Schema的约束。
解决方案
如果需要让id列的nullable显示为false,可以在加载数据后通过以下方式实现:
方案1:重新构建Schema并转换DataFrame
# 定义明确nullable=false的新Schema new_schema = StructType([ StructField("id", StringType(), False), StructField("course_id", StringType(), True), StructField("rate", DoubleType(), True), StructField("date", TimestampType(), True), StructField("display_name", StringType(), True), StructField("comment", StringType(), True), StructField("new_id", StringType(), True) ]) # 基于原DataFrame的RDD重新创建DataFrame应用新Schema comments_df_fixed = spark.createDataFrame(comments_df.rdd, schema=new_schema) comments_df_fixed.printSchema()
方案2:使用select时指定列的非空属性
from pyspark.sql.functions import col comments_df_fixed = comments_df.select( col("id").alias("id", nullable=False), col("course_id"), col("rate"), col("date"), col("display_name"), col("comment"), col("new_id") ) comments_df_fixed.printSchema()
额外校验:确保数据无空值
如果需要确保id列确实不存在空值,可以在加载后添加断言校验:
assert comments_df.filter(col("id").isNull()).count() == 0, "id列存在空值"
内容的提问来源于stack exchange,提问作者Harbeer Kadian
相关产品推荐
相关产品推荐

