PySpark写入SQL Server报列nullable配置不一致错误排查
问题
执行代码时抛出异常,提示DataFrame与SQL Server目标表中Hiring_Date列的nullable配置存在差异,报错触发位置为df3.write行,该行逻辑为将处理完成的df3写入SQL Server目标表。
错误信息
java.sql.SQLException: Spark Dataframe and SQL Server table have differing column nullable configurations at column index 0 DF col Hiring_Date nullable config is true Table col Hiring_Date nullable config is false
错误说明:Spark DataFrame与SQL Server表在索引为0的列上nullable配置不匹配,DataFrame中Hiring_Date列nullable属性为true,目标SQL表中Hiring_Date列nullable属性为false
备注
根据实际使用经验(相同处理逻辑在其他脚本中可正常运行),在PySpark中定义列数据类型并将列中空值替换为指定值后,DataFrame中对应列应为非nullable属性。
业务代码
df = spark.read.csv("myDataFile.txt", sep="|", header="true", inferSchema="false") df1 = df.select( *[ F.when(F.col(column).isNull(),'').otherwise(F.col(column)).alias(column) for column in df.columns]) df2 = df1.withColumn("Hiring_Date", df1.Hiring_Date.cast(TimestampType())) \ .withColumn("Hiring_Fee", df1.Hiring_Fee.cast(DoubleType())) df3 = df2.fillna( {'Hiring_Fee' : 0, 'Hiring_Date': '1753-01-01 00:00:00.000'} ) try: df3.write \ .format("com.microsoft.sqlserver.jdbc.spark") \ .mode("append") \ .option("url", url) \ .option("dbtable", table_name) \ .option("user", myUserName) \ .option("password", myPassword) \ .save() except ValueError as error : print("Connector write failed", error)
SQL Server目标表定义
CREATE TABLE HR_History( Hiring_Date datetime NOT NULL, Hiring_Fee float NOT NULL )
根因与遗漏配置
- 核心遗漏:完成空值填充后,没有显式修改DataFrame对应列的nullable元数据为false,和SQL Server表的NOT NULL约束对齐。
- 认知误区说明:PySpark中
when/otherwise空值替换、cast类型转换、fillna空值填充这类算子,都不会自动修改列schema中的nullable标记。这个标记是从数据源读取时静态继承的:本次代码读取CSV时设置inferSchema="false",所有列默认都是nullable=true的字符串类型,后续所有转换算子都保留了这个nullable=true的元数据,哪怕实际数据已经没有空值,元数据标记也不会自动更新。 - 其他脚本可正常运行的原因:大概率是用了不做nullable严格校验的旧版JDBC连接器,或者在写入前已经做过显式schema定义。
修复方式
在写入前,显式将目标列的nullable属性调整为false即可,参考代码:
from pyspark.sql.types import StructType # 基于现有schema调整非空配置 fixed_schema = StructType([ field.copy(nullable=False) if field.name in ("Hiring_Date", "Hiring_Fee") else field for field in df3.schema ]) # 用修正后的schema重新生成DataFrame df3 = spark.createDataFrame(df3.rdd, fixed_schema) # 后续写入逻辑保持不变即可
内容的提问来源于stack exchange,提问作者nam
相关产品推荐
相关产品推荐

