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

请求协助编写Databricks中Delta Live Tables增量数据质量校验PySpark代码

Delta Live Tables 增量同步与数据质量校验实现

1. Bronze层:增量数据同步

使用Auto Loader增量读取源数据(以CSV格式为例,可根据实际数据源调整),写入Bronze层Delta表,保留原始数据及元数据。

import dlt
from pyspark.sql.functions import current_timestamp, input_file_name

@dlt.table(
    name="bronze_raw_data",
    comment="原始增量数据的Bronze层表,保留所有原始字段及加载元数据",
    table_properties={
        "quality": "bronze",
        "delta.autoOptimize.optimizeWrite": "true",
        "delta.autoOptimize.autoCompact": "true"
    }
)
def bronze_raw_data():
    # 配置Auto Loader读取源数据,可替换为S3、GCS或其他数据源路径
    return (
        spark.readStream.format("cloudFiles")
        .option("cloudFiles.format", "csv")
        .option("cloudFiles.schemaLocation", "/dbfs/dlt/schemas/bronze_raw_data")  # 自动推断并保存Schema
        .option("cloudFiles.inferColumnTypes", "true")
        .load("/mnt/raw-source-data/")
        .withColumn("load_timestamp", current_timestamp())  # 添加数据加载时间戳
        .withColumn("source_file", input_file_name())  # 添加源文件路径
    )

2. 数据质量校验规则

通过DLT内置约束对Bronze层数据校验,筛选合格数据流向Silver层。这里用@dlt.expect_or_drop过滤违规数据,也可根据需求改用@dlt.expect_or_fail终止流水线或@dlt.expect仅记录违规信息。

# 定义校验后的视图,作为Silver层的数据源
@dlt.view(name="validated_bronze_data", comment="经过数据质量校验的Bronze层数据视图")
@dlt.expect_or_drop("non_null_id", "id IS NOT NULL")  # 过滤id为空的数据
@dlt.expect_or_drop("valid_status", "status IN ('active', 'inactive')")  # 过滤状态不在允许范围的数据
@dlt.expect_or_drop("positive_amount", "amount > 0")  # 过滤金额非正的数据
def validated_bronze_data():
    return dlt.read("bronze_raw_data")

3. Silver层:加载合格数据

从校验视图读取数据,完成清洗转换后写入Silver层Delta表,支持增量更新。

@dlt.table(
    name="silver_cleaned_data",
    comment="经过数据质量校验与清洗的Silver层表",
    table_properties={
        "quality": "silver",
        "delta.autoOptimize.optimizeWrite": "true",
        "delta.autoOptimize.autoCompact": "true"
    }
)
def silver_cleaned_data():
    return (
        dlt.read_stream("validated_bronze_data")
        .withColumnRenamed("old_column", "new_column")  # 根据实际需求重命名字段
        .select(
            "id",
            "user_name",
            "status",
            "amount",
            "load_timestamp",
            "source_file"
        )  # 选择需要保留的字段
    )

关键说明

  • Auto Loader:自动识别新增文件,支持Schema自动推断与演化,适配增量同步场景。
  • 数据质量控制:通过@dlt.expect_*系列装饰器灵活定义规则,适配不同业务容错需求。
  • Delta优化:开启autoOptimize实现自动写入优化与数据压缩,提升查询性能。

内容的提问来源于stack exchange,提问作者Venkatesh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 14:20:09