请求协助编写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
相关产品推荐
相关产品推荐

