在 Delta Live Tables 中,数据质量期望并不是只用来生成质量报告的可选插件,而是会真实参与表物化过程:一条期望可以决定违规行是继续保留、被丢弃,还是直接让管道失败。当源表只有五六列时,逐条写约束并不麻烦;但当列数增长到几十甚至上百,只约束主键、金额和状态字段远远不够,任何一列的空值或异常都可能在后续聚合分析中被无限放大。要解决这个问题,比较务实的做法是把所有列理解为一个可配置集合,通过遍历 schema 自动生成基础质量规则,再对特殊列做覆盖和排除。

一、为什么全列质量约束比逐条手写更可靠
手写数据质量规则最常见的失败模式,并不是约束逻辑写错,而是覆盖范围不够。一个订单表可能有三十个字段,开发人员通常会把注意力放在订单号、金额、状态和日期上,而客户备注、渠道标识、设备号等字段往往只在出现问题时才补规则。批量生成列级约束可以强制让每一列至少经过基础检查,例如非空、字符串长度上限或枚举值范围。
Delta Live Tables 的期望机制有一个明显优势:约束定义紧邻表定义,并且被记录到事件日志中。即使某条规则只是保留违规行而不中断管道,你也能在后续监控中看到每一列的失败次数。批量生成规则时,建议把规则名设计得稳定且可读,例如 not_null_customer_id 或 length_phone,而不是直接用自增编号。这样当质量报表出现异常时,定位列会非常快。
二、用Python动态生成列级期望的三种写法
最直接的方式是在装饰器外读取源表 schema,生成一个规则字典,再交给 @dlt.expect_all。这种方式的优点是所有列都经过同一套逻辑,适合以非空检查为主的场景。你可以维护一个排除集合,把允许为空的备注列或已经单独约束的列跳过。
import dlt
EXCLUDE_COLUMNS = {"comment", "raw_payload"}
def build_not_null_rules(columns, exclude=None):
exclude = exclude or set()
rules = {}
for column in columns:
if column in exclude:
continue
rule_name = f"not_null_{column}"
expression = f"`{column}` IS NOT NULL"
rules[rule_name] = expression
return rules
@dlt.table
@dlt.expect_all(build_not_null_rules(spark.read.table("orders_bronze").columns, EXCLUDE_COLUMNS))
def orders_silver():
return dlt.read("orders_bronze")
这段代码在创建表之前先读取 spark.read.table("orders_bronze").columns,因此导入阶段就确定了规则集合。需要注意的是,如果上游表结构频繁变化,例如新增列或删除列,规则集会随之自动调整,这是预期行为;但如果某列被误删,对应的约束也会静默消失。因此最好把排除集合放进版本控制,而不是散落在 notebook 里。
第二种写法是把循环放进表函数内部,一边读取数据,一边对 DataFrame 应用 dlt.expect。这种方式更灵活,因为可以先做一些列清理或重命名,再基于实际列名生成规则。
@dlt.table
def orders_silver():
df = dlt.read("orders_bronze")
for column in df.columns:
if column in ("comment", "raw_payload"):
continue
df = dlt.expect(
f"not_null_{column}",
f"`{column}` IS NOT NULL"
)(df)
return df
第三种写法是引入元数据驱动规则。把每个列名、规则类型、参数放到一张配置表或字典中,再统一解析。例如 phone 列执行长度检查,amount 列执行非空和正数检查。这样不必为每种规则单独写分支,也能让非技术用户通过配置表调整质量策略。下面是一个简化示例。
column_rules = {
"customer_id": ["not_null"],
"phone": ["not_null", "length:8:20"],
"amount": ["not_null", "positive"]
}
def build_rules_from_metadata(metadata):
rules = {}
for column, rule_list in metadata.items():
for rule in rule_list:
if rule == "not_null":
rules[f"not_null_{column}"] = f"`{column}` IS NOT NULL"
elif rule == "positive":
rules[f"positive_{column}"] = f"`{column}` > 0"
elif rule.startswith("length:"):
_, min_len, max_len = rule.split(":")
rules[f"length_{column}"] = f"length(cast(`{column}` as string)) BETWEEN {min_len} AND {max_len}"
return rules
@dlt.table
@dlt.expect_all(build_rules_from_metadata(column_rules))
def orders_silver():
return dlt.read("orders_bronze")
三、SQL声明式实现与混合模式
如果团队主要使用 SQL 维护 Delta Live Tables,也可以使用 CONSTRAINT ... EXPECT 语法。静态写法非常直观,尤其适合列数不多但规则明确的表。
CREATE OR REFRESH STREAMING LIVE TABLE orders_silver AS SELECT * FROM STREAM(LIVE.orders_bronze) CONSTRAINT not_null_id EXPECT (`id` IS NOT NULL) ON VIOLATION DROP ROW CONSTRAINT not_null_order_date EXPECT (`order_date` IS NOT NULL) ON VIOLATION DROP ROW CONSTRAINT valid_amount EXPECT (`amount` IS NOT NULL AND `amount` > 0) ON VIOLATION DROP ROW
当列数很多时,手动在 SQL 中列出每一种约束会比较啰嗦,也容易因为列名改动而失效。此时可以让 Python 读取 schema,自动生成 SQL 片段。生成逻辑和 Python 装饰器方式一样,只是最终输出的是 DDL 字符串。
columns = spark.read.table("orders_bronze").columns
exclude = {"comment", "raw_payload"}
constraints = ",\n".join([
f"CONSTRAINT not_null_{c} EXPECT (`{c}` IS NOT NULL) ON VIOLATION DROP ROW"
for c in columns if c not in exclude
])
sql = f"""
CREATE OR REFRESH STREAMING LIVE TABLE orders_silver
AS SELECT * FROM STREAM(LIVE.orders_bronze)
{constraints}
"""
print(sql)
混合模式的好处是可以复用组织内部已有的 Python 规则库,同时让表定义仍然保留在 SQL 中,方便 DBA 和数据分析师阅读。但要注意,SQL 字符串一旦被动态生成,执行前最好打印或记录完整 DDL,否则排查管道问题时很难还原实际约束内容。
四、失败策略、性能开销与监控
批量应用约束时,必须根据业务要求选择失败策略。Delta Live Tables 常见的三类操作分别是保留并记录、丢弃违规行、直接失败。保留并记录适合探索性数据,丢弃违规行适合下游严格要求的宽表,直接失败适合核心事实表的硬性业务规则。不同装饰器对应不同行为。
| 策略 | 常用写法 | 适用场景 |
|---|---|---|
| 保留并记录 | @dlt.expect_all | 允许坏行进入表,但需要统计失败比例 |
| 丢弃违规行 | @dlt.expect_all_or_drop | 下游无法容忍空值或格式错误的宽表 |
| 中断管道 | @dlt.expect_all_or_fail | 核心字段出现严重问题时必须停止任务 |
性能方面,列级非空、类型检查表达式的开销通常很低,Spark Catalyst 会把多个简单条件合并到同一次扫描中。真正需要关注的是复杂正则、JSON 解析或自定义 UDF,它们会阻止谓词下推,甚至导致数据倾斜。建议把基础质量规则控制在 Spark SQL 内建函数范围内,尽量避免 Python UDF 出现在约束里。
最后要建立质量监控闭环。每次管道运行后,应查询事件日志中的期望指标,观察失败行数是否突然上升。比如某个第三方数据源开始把手机号写入空字符串,行数没变但长度规则失败率会立刻提高。把失败指标接入告警后,批量列级约束才能真正从静态规则变成可运维的质量防线。
Delta Live Tables数据质量期望约束修改时间:2026-09-17 22:54:08