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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 03:40:05