数据质量为什么是 ETL 的生命线

先讲一个真实的故事。某电商公司每天凌晨跑 ETL,把前一天的订单数据导入分析系统。数据团队每天早上的第一件事就是检查报表。有一天,市场部开会时发现 GMV(总成交额)比前一天跌了 80%。大家慌了——不是运营出了问题,而是 ETL 任务在凌晨 2 点因为源库连接超时失败了,但没有人发现。整个上午的决策都基于错误的数据。

这种情况在数据工程领域太常见了。数据管道的输出质量直接影响着下游的报表、分析和机器学习模型。没有质量监控的 ETL 管道,就像没有质检的生产线——你不知道产品里有没有混入次品,甚至不知道生产线是不是还在运转。

数据质量的六个维度

在建立数据质量监控之前,先要明确"质量好"到底是什么意思。业界通常用这六个维度来衡量:

准确性(Accuracy)

数据是否真实反映了现实世界的情况?一个用户的年龄写的是 200 岁,这不准确。一笔订单金额是负数,这也是不准确。

完整性(Completeness)

是否有缺失的数据?订单表里必须有订单号,如果订单号字段为空,就违反了完整性要求。在 ETL 场景中,完整性还指目标表的数据行数是否和源表一致。

一致性(Consistency)

数据在不同系统之间是否一致?用户表里姓名字段在数据库 A 中是 UTF-8 编码,在数据库 B 中是 GBK 编码,同一用户在两个系统中的姓名可能不一致。

时效性(Timeliness)

数据是否在预期的时间内可用?每日报表要求在早上 8 点前准备好,如果 ETL 任务跑到了中午 12 点,数据就是"过期"的。

唯一性(Uniqueness)

数据是否有重复?主键是否唯一?在 ETL 过程中,如果同一个订单被抽取了两次,就会破坏唯一性约束。

有效性(Validity)

数据是否满足定义的格式和规则?邮件地址必须包含 @,手机号必须是 11 位数字,枚举字段必须在定义的值范围内。

在 ETL 管道中嵌入质量检查

质量检查不应该在数据入仓之后再去做,而应该嵌入在 ETL 管道的每一个环节。

抽取阶段:源端检查

在抽取数据之前,先对源数据做基本检查:

def validate_source(source_config):
    """源端数据验证"""
    checks = {
        "table_exists": check_table_exists(source_config),
        "row_count_reasonable": check_row_count(source_config),
        "freshness": check_data_freshness(source_config),
        "schema_match": check_schema(source_config),
    }
    
    failed = [name for name, passed in checks.items() if not passed]
    if failed:
        raise DataQualityError(f"源端质量检查失败: {failed}")

转换阶段:行级检查

在数据转换过程中,对每一条数据进行验证:

import pandas as pd
import re


def validate_record(row):
    """单行数据质量校验"""
    errors = []
    
    # 1. 非空检查
    required_fields = ["order_id", "user_id", "amount"]
    for field in required_fields:
        if pd.isna(row.get(field)):
            errors.append(f"{field} 为空")
    
    # 2. 格式检查
    email_pattern = r'^[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}$'
    if not re.match(email_pattern, str(row.get("email", ""))):
        errors.append("邮箱格式无效")
    
    # 3. 范围检查
    amount = row.get("amount", 0)
    if amount < 0:
        errors.append("金额不能为负数")
    elif amount > 1000000:
        errors.append("金额超出正常范围")
    
    # 4. 枚举值检查
    valid_statuses = ["PENDING", "PAID", "SHIPPED", "DELIVERED", "CANCELLED"]
    if row.get("status") not in valid_statuses:
        errors.append(f"订单状态无效: {row.get('status')}")
    
    return errors


def transform_with_quality(df):
    """带质量检查的转换"""
    valid_rows = []
    bad_rows = []
    
    for idx, row in df.iterrows():
        errors = validate_record(row)
        if errors:
            bad_rows.append({
                "index": idx,
                "record": row.to_dict(),
                "errors": errors
            })
        else:
            valid_rows.append(row)
    
    # 记录质量报告
    quality_report = {
        "total": len(df),
        "valid": len(valid_rows),
        "invalid": len(bad_rows),
        "error_rate": len(bad_rows) / len(df) * 100 if len(df) > 0 else 0,
        "sample_errors": bad_rows[:10]
    }
    
    return pd.DataFrame(valid_rows), quality_report

加载阶段:完整性检查

数据写入目标后,做总量级的比对验证:

-- 行数比对检查
SELECT 
    'source' AS source_type,
    COUNT(*) AS row_count,
    COUNT(DISTINCT order_id) AS unique_count
FROM staging.orders
WHERE dt = CURRENT_DATE

UNION ALL

SELECT 
    'target' AS source_type,
    COUNT(*) AS row_count,
    COUNT(DISTINCT order_id) AS unique_count
FROM dw.fact_orders
WHERE dt = CURRENT_DATE;

聚合级质量指标

除了行级别的检查,还需要关注数据整体的统计特征是否合理。这些"听起来不对劲"的情况往往能暴露更深层次的问题。

def compute_quality_metrics(df, table_name):
    """计算聚合质量指标"""
    metrics = {}
    
    # 1. 统计特征
    metrics["row_count"] = len(df)
    metrics["null_rates"] = df.isnull().mean().to_dict()
    metrics["duplicate_rate"] = df.duplicated().mean()
    
    # 2. 数值分布
    numeric_cols = df.select_dtypes(include=["number"]).columns
    for col in numeric_cols:
        metrics[f"{col}_mean"] = df[col].mean()
        metrics[f"{col}_std"] = df[col].std()
        metrics[f"{col}_min"] = df[col].min()
        metrics[f"{col}_max"] = df[col].max()
        metrics[f"{col}_p99"] = df[col].quantile(0.99)
    
    # 3. 分类分布
    cat_cols = df.select_dtypes(include=["object"]).columns
    for col in cat_cols[:5]:  # 只检查前5个分类列
        value_counts = df[col].value_counts(normalize=True)
        metrics[f"{col}_top_value"] = value_counts.index[0]
        metrics[f"{col}_top_ratio"] = value_counts.iloc[0]
        metrics[f"{col}_unique_count"] = df[col].nunique()
    
    return metrics


def detect_anomaly(current_metrics, historical_metrics, threshold=3):
    """基于历史数据检测异常"""
    anomalies = []
    
    for key, current_value in current_metrics.items():
        if key not in historical_metrics:
            continue
            
        historical = historical_metrics[key]
        mean = historical["mean"]
        std = historical["std"]
        
        if std == 0:
            continue
            
        z_score = abs(current_value - mean) / std
        if z_score > threshold:
            anomalies.append({
                "metric": key,
                "current": current_value,
                "expected": mean,
                "z_score": z_score,
                "severity": "CRITICAL" if z_score > 5 else "WARNING"
            })
    
    return anomalies

异常检测策略

聚合指标只能告诉你"出问题了",但不会告诉你"哪里出问题了"。更主动的异常检测策略包括以下几种。

基线对比

维护一个历史基线,每次 ETL 完成后自动对比当前数据和基线的差异。比如今天导入了 95 万条订单,但过去一周平均是 100 万条,相差 5%。如果偏差超过设定的阈值(比如 10%),就触发告警。

规则引擎

基于业务规则进行验证。这类规则需要业务和数据团队共同制定:

class DataQualityRule:
    """数据质量规则基类"""
    def check(self, df) -> bool:
        raise NotImplementedError


class CompletenessRule(DataQualityRule):
    """完整性规则:指定字段不可为空"""
    def __init__(self, columns, threshold=0.99):
        self.columns = columns
        self.threshold = threshold
    
    def check(self, df):
        for col in self.columns:
            completeness = 1 - df[col].isnull().mean()
            if completeness < self.threshold:
                print(f"字段 {col} 的完整度 {completeness:.2%} 低于阈值 {self.threshold:.2%}")
                return False
        return True


class UniquenessRule(DataQualityRule):
    """唯一性规则:指定组合必须唯一"""
    def __init__(self, columns):
        self.columns = columns
    
    def check(self, df):
        dup_count = df.duplicated(subset=self.columns).sum()
        if dup_count > 0:
            print(f"发现 {dup_count} 条重复数据,字段: {self.columns}")
            return False
        return True


class ReferentialIntegrityRule(DataQualityRule):
    """引用完整性规则:外键必须在主表中存在"""
    def __init__(self, fk_column, reference_df, pk_column):
        self.fk_column = fk_column
        self.reference_df = reference_df
        self.pk_column = pk_column
    
    def check(self, df):
        valid_ids = set(self.reference_df[self.pk_column])
        invalid_mask = ~df[self.fk_column].isin(valid_ids)
        invalid_count = invalid_mask.sum()
        
        if invalid_count > 0:
            print(f"发现 {invalid_count} 条无效的外键引用")
            return False
        return True


# 使用质量规则引擎
quality_engine = [
    CompletenessRule(["order_id", "user_id", "amount"]),
    UniquenessRule(["order_id"]),
    ReferentialIntegrityRule("user_id", users_df, "user_id"),
]

all_passed = all(rule.check(orders_df) for rule in quality_engine)

异常处理策略

即使有了最好的质量检查,异常数据仍然会出现。关键不在于杜绝所有异常,而在于异常发生时系统能正确处理。

拒绝(Reject)

不符合质量要求的数据直接拒绝入库,记录到错误日志中供人工排查。

隔离(Quarantine)

把异常数据放入独立的隔离区,不影响正常数据的可用性。

def quarantine_bad_data(bad_records, quarantine_table="data_quality_quarantine"):
    """将异常数据隔离到专门区域"""
    for record in bad_records:
        quarantine_sql = f"""
        INSERT INTO {quarantine_table} (
            source_table, record_id, raw_data, error_type,
            error_message, created_at, status
        ) VALUES (%s, %s, %s, %s, %s, NOW(), 'PENDING_REVIEW')
        """
        execute(quarantine_sql, (
            record["source_table"],
            record["record_id"],
            record["raw_data"],
            record["error_type"],
            record["error_message"],
        ))

转换(Transform)

如果异常数据可以通过修复逻辑自动处理,就在管道中加上转换逻辑:

def auto_fix_bad_data(df):
    """自动修复常见数据问题"""
    
    # 修复空字符串
    for col in df.columns:
        df[col] = df[col].replace("", None)
    
    # 修复常见日期格式
    def parse_date_flexible(value):
        if pd.isna(value):
            return None
        for fmt in ["%Y-%m-%d", "%Y/%m/%d", "%d-%m-%Y", "%Y%m%d"]:
            try:
                return pd.to_datetime(value, format=fmt)
            except (ValueError, TypeError):
                continue
        return None
    
    if "order_date" in df.columns:
        df["order_date"] = df["order_date"].apply(parse_date_flexible)
    
    # 修复负数金额
    if "amount" in df.columns:
        df["amount"] = df["amount"].abs()
    
    # 去除首尾空格
    string_cols = df.select_dtypes(include=["object"]).columns
    for col in string_cols:
        df[col] = df[col].str.strip()
    
    return df

Great Expectations 实践

Great Expectations 是目前最流行的开源数据质量工具。它把数据质量检查从"写 Python 函数"升级为"声明期望 + 自动生成文档"。

安装和基本使用

pip install great_expectations

定义期望

import great_expectations as ge

# 加载数据
df = ge.read_csv("orders.csv")

# 定义期望
expectation_suite = {
    "expect_column_values_to_not_be_null": {
        "column": "order_id"
    },
    "expect_column_values_to_be_unique": {
        "column": "order_id"
    },
    "expect_column_values_to_be_between": {
        "column": "amount",
        "min_value": 0,
        "max_value": 1000000
    },
    "expect_column_values_to_be_in_set": {
        "column": "status",
        "value_set": ["PENDING", "PAID", "SHIPPED", "DELIVERED", "CANCELLED"]
    },
    "expect_column_value_lengths_to_be_between": {
        "column": "phone",
        "min_value": 10,
        "max_value": 15
    },
    "expect_column_mean_to_be_between": {
        "column": "amount",
        "min_value": 50,
        "max_value": 500
    }
}

# 逐个验证
results = []
for expectation_name, kwargs in expectation_suite.items():
    method = getattr(df, expectation_name)
    result = method(**kwargs)
    results.append({
        "expectation": expectation_name,
        "column": kwargs.get("column"),
        "success": result["success"],
        "unexpected_count": result["result"].get("unexpected_count", 0),
    })

# 打印结果
for r in results:
    status = "✅" if r["success"] else "❌"
    print(f"{status} {r['expectation']}: {r['unexpected_count']} 条异常")

在 ETL 管道中集成

def etl_with_quality_gate(source_path, target_table):
    """带质量门的 ETL 流程"""
    # 1. 抽取
    df = ge.read_csv(source_path)
    
    # 2. 质量门:检查通过才继续
    quality_check = df.expect_column_values_to_not_be_null("order_id")
    if not quality_check["success"]:
        raise DataQualityError("质量门拒绝: order_id 存在空值")
    
    # 3. 转换
    df = transform(df)
    
    # 4. 质量门:检查转换后数据
    amount_check = df.expect_column_values_to_be_between(
        "amount", min_value=0, max_value=1000000
    )
    if not amount_check["success"]:
        quarantine_bad_data(df[df["amount"] < 0], target_table)
        df = df[df["amount"] >= 0]
    
    # 5. 加载
    df.to_sql(target_table, conn, if_exists="append", index=False)

告警与通知策略

质量检查做完了,发现有问题。然后呢?光有检查没有告警等于白做。一个好的告警策略要考虑三个要素:

分级告警

不是所有问题都需要马上叫醒值班工程师:

def alert_if_needed(quality_report):
    """分级告警"""
    error_rate = quality_report["error_rate"]
    
    if error_rate > 30:
        # P0: 严重问题,30% 以上数据异常
        send_phone_alert("ETL 管道严重质量问题,错误率 {error_rate:.1f}%")
        create_incident_ticket("critical", quality_report)
    elif error_rate > 10:
        # P1: 重要问题,需要关注
        send_slack_alert(f"ETL 质量问题告警,错误率 {error_rate:.1f}%")
    elif error_rate > 1:
        # P2: 一般问题,记录即可
        log_to_dashboard("quality", quality_report)
    else:
        # 正常范围
        log_to_dashboard("quality", quality_report)

告警去重

同一问题连续告警只会让人麻木。实现告警聚合和静默机制:

def dedup_alert(alert_key, cooldown_minutes=30):
    """告警去重"""
    last_alert_at = cache.get(f"alert:{alert_key}")
    if last_alert_at:
        elapsed = (datetime.now() - last_alert_at).total_seconds()
        if elapsed < cooldown_minutes * 60:
            return False  # 冷却期内,不重复告警
    cache.set(f"alert:{alert_key}", datetime.now(), ttl=3600)
    return True

建立数据质量仪表盘

质量检查的最终产出应该是一个可视化的仪表盘,让每个人都能清楚地看到数据管道的健康状况。关键指标包括:

  • 数据新鲜度:数据最后更新的时间
  • 管道成功率:最近 N 次 ETL 任务的成功率
  • 质量通过率:通过质量检查的数据比例
  • 异常趋势:质量问题的数量随时间的变化
  • 修复时效:从发现异常到修复的平均时间

小结

数据质量不是 ETL 管道的附加功能,而是核心组成部分。在管道中嵌入质量检查,建立异常发现和处理的完整流程,是数据工程成熟度的重要标志。从行级别检查到聚合指标对比,从简单拒绝到智能修复,质量监控体系的建设是一个持续迭代的过程。选择 Great Expectations 这样的工具可以帮助团队更标准化地管理质量检查,但更重要的是建立"数据必有质量门"的工程文化。

Summary: 数据质量六大维度与监控体系,从行级检查到聚合异常检测,集成质量门和告警策略。