You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 09:23:20