Azure Databricks中PySpark读取JSONL文件Schema校验失效问题
问题分析
当前使用PERMISSIVE模式读取JSONL时,仅会捕获JSON语法完全错误的记录(如非合法JSON结构),对于以下场景不会标记为脏数据:
- 非空字段为
null - 字段缺失
- 字段类型不匹配(解析后对应字段设为
null)
Schema中的nullable=False仅作为元数据存在,Spark读取阶段不会强制校验该约束,仅在写入Delta表启用约束校验时才会触发报错。
解决方案
通过保留原始JSON记录+自定义校验逻辑的方式,精准捕获所有不符合要求的脏数据:
步骤1:读取原始JSON字符串
先将每条JSON记录以字符串形式读取,保留原始内容用于后续标记脏数据:
# 读取原始JSONL文件,每条记录作为字符串存入original_record字段 raw_df = spark.read.text(raw_file_location).withColumnRenamed("value", "original_record")
步骤2:解析JSON并校验脏数据
使用from_json解析到指定Schema,同时添加自定义校验条件,将不符合要求的记录原始内容存入_corrupt_record:
from pyspark.sql.functions import from_json, col, when # 解析JSON到预定义Schema parsed_df = raw_df.withColumn("parsed_data", from_json(col("original_record"), schema)) # 提取字段并标记脏数据 final_df = parsed_df.select( col("parsed_data.restaurantId"), col("parsed_data.reviewId"), col("parsed_data.text"), col("parsed_data.rating"), col("parsed_data.publishedAt"), # 校验所有非空字段是否为null,或解析失败(parsed_data为null) when( col("parsed_data").isNull() | col("parsed_data.restaurantId").isNull() | col("parsed_data.reviewId").isNull() | col("parsed_data.text").isNull() | col("parsed_data.rating").isNull() | col("parsed_data.publishedAt").isNull(), col("original_record") ).alias("_corrupt_record") )
步骤3:验证脏数据
运行以下代码查看捕获到的脏数据:
display(final_df.filter(col("_corrupt_record").isNotNull()))
内容的提问来源于stack exchange,提问作者prabhudotpy
相关产品推荐
相关产品推荐

