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

如何为DLT表中所有列应用数据质量校验规则?

问题描述

查阅了各类教程与文章,发现大多介绍的是为DLT表中部分列配置数据质量校验规则,示例代码如下:

@dlt.table(
  comment="Wikipedia clickstream data cleaned and prepared for analysis."
)
@dlt.expect("valid_current_page_title", "current_page_title IS NOT NULL")
@dlt.expect_or_fail("valid_count", "click_count > 0")
def clickstream_prepared():
  return (
    dlt.read("clickstream_raw")
      .withColumn("click_count", expr("CAST(n AS INT)"))
      .withColumnRenamed("curr_title", "current_page_title")
      .withColumnRenamed("prev_title", "previous_page_title")
      .select("current_page_title", "click_count", "previous_page_title")
  )

此处需手动指定要校验的列,但我希望为DataFrame中的所有列应用校验规则。我曾尝试通过循环动态生成函数来实现,但该方法效率极低,代码如下:

for column in columns_list_order_table:     
    exec(f'''
@dlt.table(comment="null value validations for {column}")
@dlt.expect_or_drop("null values","is_null == false")
def null_validation_orders_for_column_{column}():
    df = dlt.read("bronze_orders")
    return df.withColumn("is_null", col("{column}").isNull())
''')

请问如何高效实现为DLT表所有列应用数据质量校验?

高效实现方案

方法1:使用@dlt.expect_all系列装饰器批量定义规则

DLT提供了expect_all、expect_all_or_drop、expect_all_or_fail装饰器,支持一次性传入包含所有列校验规则的字典,无需逐个添加装饰器。以校验所有列非空为例:

@dlt.table(comment="Bronze orders with full column null validation")
@dlt.expect_all_or_drop({
    f"valid_{col}": f"{col} IS NOT NULL" 
    for col in dlt.read("bronze_orders").columns
})
def silver_orders_validated():
    return dlt.read("bronze_orders")

只需替换expect_all_or_drop为对应装饰器,即可切换校验失败后的处理逻辑(记录告警/丢弃行/终止任务)。

方法2:DataFrame层面动态生成校验逻辑

如果需要自定义不同列的校验规则(比如部分列校验非空,部分列校验数值范围),可以在DataFrame中动态生成并合并校验条件:

from pyspark.sql.functions import col

@dlt.table(comment="Bronze orders with dynamic column validation")
def silver_orders_dynamic():
    df = dlt.read("bronze_orders")
    
    # 生成所有列的校验条件,这里以非空为例,可按需修改
    validation_conditions = []
    for col_name in df.columns:
        # 示例:对数值列额外校验大于0
        if df.schema[col_name].dataType.simpleString() in ["int", "bigint", "double"]:
            validation_conditions.append(col(col_name).isNotNull() & (col(col_name) > 0))
        else:
            validation_conditions.append(col(col_name).isNotNull())
    
    # 合并所有条件,保留符合所有规则的行
    combined_condition = validation_conditions[0]
    for cond in validation_conditions[1:]:
        combined_condition = combined_condition & cond
    
    return df.filter(combined_condition)

这种方式灵活性更高,适合复杂的多规则校验场景。

方法3:动态添加装饰器(替代exec)

如果偏好使用装饰器模式但不想用exec执行动态代码,可以通过Python的装饰器API为函数批量添加校验规则:

from functools import wraps

def add_full_column_validation(table_name, validation_type="expect"):
    def decorator(func):
        # 获取目标表的列
        df = dlt.read(table_name)
        # 选择对应的DLT校验装饰器
        dlt_decorator = {
            "expect": dlt.expect,
            "expect_or_drop": dlt.expect_or_drop,
            "expect_or_fail": dlt.expect_or_fail
        }[validation_type]
        
        # 批量添加校验装饰器
        wrapped_func = func
        for col_name in df.columns:
            wrapped_func = dlt_decorator(f"valid_{col_name}", f"{col_name} IS NOT NULL")(wrapped_func)
        
        # 添加table装饰器并返回
        return dlt.table(comment=f"Full validation for {table_name}")(wrapped_func)
    return decorator

# 使用自定义装饰器
@add_full_column_validation(table_name="bronze_orders", validation_type="expect_or_fail")
def silver_orders_strict_validated():
    return dlt.read("bronze_orders")

这种方式比exec更安全、易维护,避免了动态代码执行的风险。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 00:34:59