Databricks Spark中Schema设为nullable=false仍出现空值的问题
问题排查:Spark读取CSV时非空字段出现Null且未写入坏记录路径
背景信息
定义的JSON Schema
json_schema_string = "{\"fields\":[{\"metadata\":{},\"name\":\"Name\",\"nullable\":true,\"type\":\"string\"},{\"metadata\":{},\"name\":\"empid\",\"nullable\":false,\"type\":\"integer\"},{\"metadata\":{},\"name\":\"age\",\"nullable\":false,\"type\":\"integer\"},{\"metadata\":{},\"name\":\"ph_number\",\"nullable\":true,\"type\":\"long\"},{\"metadata\":{},\"name\":\"address\",\"nullable\":true,\"type\":\"string\"}],\"type\":\"struct\"}"
转换后的StructType
StructType([StructField('Name', StringType(), True), StructField('empid', IntegerType(), False), StructField('age', IntegerType(), False), StructField('ph_number', LongType(), True), StructField('address', StringType(), True)])
执行的Spark读取CSV代码
df_current = spark.read\ .option("header","true")\ .schema(json_schema)\ .option("inferSchema","false")\ .option("badRecordsPath", f"{bad_records_path_session}/{STEP_NAME}") \ .option("multiLine", "true")\ .option('escape', "\"")\ .option('delimiter', ',')\ .option("encoding", "UTF-8")\ .csv(file_path)\ .selectExpr("*", "lower(_metadata.file_name) as pmd_file_path", "_metadata.file_modification_time as pmd_file_mod_DTM", f"{time_of_run} as pmd_batch_id") ##df_current.write.format("delta").mode("append").saveAsTable(f"{catalog_name}.dnu_development.test_reject_records") df_current.show()
问题现象
执行df_current.show()时发现empid、age字段存在Null值,但这些不符合Schema非空要求的行并未被写入badRecordsPath指定的拒绝文件路径中。
排查与解决方案
1. 明确badRecordsPath的触发逻辑
Spark的badRecordsPath仅在数据解析时发生类型转换错误(比如字符串无法转成整数)时才会将行写入坏记录路径。而CSV中对应字段为空字符串的情况,Spark会自动将其转为Null,这不属于解析错误,因此不会触发坏记录捕获。
针对空字符串转Null的场景,需手动处理:
- 读取时添加
option("nullValue", ""),明确将空字符串映射为Null; - 读取后拆分有效行与无效行,将不符合非空要求的行写入指定路径:
# 拆分有效/无效数据 valid_df = df_current.filter(df_current.empid.isNotNull() & df_current.age.isNotNull()) invalid_df = df_current.filter(df_current.empid.isNull() | df_current.age.isNull()) # 写入无效行到坏记录路径 invalid_df.write.format("json").mode("append").save(f"{bad_records_path_session}/{STEP_NAME}") # 处理有效数据 valid_df.show()
2. 验证Schema的实际生效情况
通过print(df_current.schema)输出DataFrame的Schema,确认empid和age字段的nullable属性确实为False,避免转换过程中出现配置丢失。
3. 检查CSV文件的空值格式
- 如果CSV中
empid/age列的空值是空白字符(空格、制表符),添加option("trim", "true")参数去除字段前后空格,此时空白字符会被转为空字符串再转Null,若需要触发解析错误,可结合步骤1的手动过滤; - 如果是空字符串,直接按步骤1的手动拆分逻辑处理。
4. 确认Spark版本兼容性
badRecordsPath在Spark 2.4及以上版本才具备完整功能,若使用旧版本,建议升级到兼容版本,或改用手动过滤的方式处理无效行。
内容的提问来源于stack exchange,提问作者Rakesh Prasad
相关产品推荐
相关产品推荐

