5 minutes
数据质量监控与异常处理
数据质量为什么是 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: 数据质量六大维度与监控体系,从行级检查到聚合异常检测,集成质量门和告警策略。