项目背景
在一家日均订单量约 5 万单的中型电商平台中,数据团队发现近三个月 BI 报表中的"月度 GMV"与财务系统对账时存在约 3.2% 的偏差。经过排查,问题根源在于数据从业务库同步到数据仓库的 ETL 过程中缺乏质量校验,导致"脏数据"进入了下游分析系统。
本项目模拟该真实业务场景,构建了一套数据质量自动化校验体系,覆盖完整性、一致性、准确性、及时性、唯一性五大质量维度。
项目目标
- 理解数据质量对业务决策的重要性
- 掌握 DAMA 数据治理框架中的五大质量维度
- 使用 Python + Pandas 实现自动化数据质量校验
- 输出标准化的数据质量测试报告
- 建立"发现问题 → 定位根因 → 修复建议"的完整闭环
业务场景与数据流
业务数据库(MySQL)
│
▼ 每日凌晨 T+1 同步
数据仓库 ODS 层(原始数据)
│
▼ 本项目的校验位置 ← 在这里拦截脏数据!
数据质量校验脚本(Python + Pandas)
│
├── 通过校验 ──→ 写入 DWD 层(干净数据,供 BI 使用)
│
└── 未通过校验 ──→ 写入异常数据表 → 触发告警 → 人工修复
数据源说明
| 数据表 | 文件名 | 格式 | 记录数 | 说明 |
|---|---|---|---|---|
| 用户表 | users.csv | CSV | 14 行 | user_id、姓名、邮箱、手机号、注册日期、年龄、城市 |
| 订单表 | orders.csv | CSV | 15 行 | 订单 ID、用户 ID、商品 ID、下单日期、金额、状态、支付方式 |
| 商品表 | products.json | JSON | 15 行 | 商品 ID、名称、分类、价格、库存 |
选型理由:
- CSV 格式:电商系统中最常见的数据交换格式,几乎所有数据库都支持 CSV 导出。
- JSON 格式:模拟 NoSQL 数据库或 REST API 返回的商品数据,体现处理多格式数据的能力。
技术栈
| 技术 | 用途 |
|---|---|
| Python 3.10+ | 主语言 |
| Pandas | 数据读取、清洗、统计分析 |
| SQLite | 结构化校验扩展 |
| logging | 日志记录 |
| openpyxl | Excel 报告导出 |
数据质量校验规则
本项目共设计 19 条校验规则,覆盖 5 大数据质量维度:
一、完整性(Completeness)
| 编号 | 规则描述 | 为什么测这个? |
|---|---|---|
| R01 | user_id、user_name、email 不能为空 | 主键缺失无法唯一标识用户;email 缺失无法发送营销邮件 |
| R02 | order_id、user_id、amount 不能为空 | order_id 是主键;user_id 缺失导致订单无法归属;amount 缺失导致 GMV 统计失真 |
| R03 | product_id、product_name、price 不能为空 | product_id 是主键;price 缺失影响定价和促销 |
| R04 | 订单状态不能为空 | 状态为空会导致订单流转中断 |
| R05 | 商品库存不能为空 | 库存为空影响"是否可购买"判断 |
二、一致性(Consistency)
| 编号 | 规则描述 | 为什么测这个? |
|---|---|---|
| R06 | 订单的 user_id 必须在用户表中存在 | 防止"孤儿订单" |
| R07 | 订单的 product_id 必须在商品表中存在 | 防止"幽灵商品" |
| R08 | 订单状态必须在合法枚举范围内 | 避免统计 SQL 漏掉异常状态 |
| R09 | 支付方式必须在合法枚举范围内 | 防止对账系统匹配失败 |
三、准确性(Accuracy)
| 编号 | 规则描述 | 为什么测这个? |
|---|---|---|
| R10 | 邮箱格式必须符合规范 | 无效邮箱会导致营销邮件全部退信 |
| R11 | 手机号格式必须为 11 位 1 开头 | 无效手机号导致短信验证码发送失败 |
| R12 | 订单金额必须大于 0 | 金额为负或为零会影响财务结算 |
| R13 | 商品价格必须 ≥ 0 | 负价格意味着"商家倒贴" |
| R14 | 用户年龄应在 1~100 之间 | 极端年龄会拉偏用户画像分析 |
四、及时性(Timeliness)
| 编号 | 规则描述 | 为什么测这个? |
|---|---|---|
| R15 | 用户注册日期不能是未来日期 | 可能是测试数据混入生产库 |
| R16 | 订单日期不能是未来日期 | 未来订单会导致实时大屏数据失真 |
五、唯一性(Uniqueness)
| 编号 | 规则描述 | 为什么测这个? |
|---|---|---|
| R17 | 用户 ID 不能重复 | 主键重复会导致 JOIN 时产生笛卡尔积 |
| R18 | 订单 ID 不能重复 | 重复订单会导致 GMV 虚高 |
| R19 | 商品 ID 不能重复 | 重复商品会导致库存扣减错乱 |
核心实现
1. 数据加载模块
def load_data(data_dir: str) -> dict:
data_dir = Path(data_dir)
data = {}
# CSV 用 dtype=str 读取,避免 ID 类字段丢失前导零
data['users'] = pd.read_csv(data_dir / "users.csv", dtype=str)
data['orders'] = pd.read_csv(data_dir / "orders.csv", dtype=str)
# JSON 格式数据
with open(data_dir / "products.json", 'r', encoding='utf-8') as f:
products_list = json.load(f)
data['products'] = pd.DataFrame(products_list, dtype=str)
return data
2. 缺陷收集器
class DefectCollector:
def __init__(self):
self.defects = []
def add(self, rule_id: str, dimension: str, table: str,
row_index: int, column: str, value, description: str,
severity: str = "中"):
self.defects.append({
'规则编号': rule_id,
'质量维度': dimension,
'表名': table,
'行号': row_index + 2, # 表头占一行,行号从 1 开始
'字段名': column,
'问题值': str(value),
'问题描述': description,
'严重等级': severity
})
def report(self) -> pd.DataFrame:
return pd.DataFrame(self.defects)
def summary(self) -> dict:
df = self.report()
if len(df) == 0:
return {'总缺陷数': 0, '通过率': '100%'}
return {
'总缺陷数': len(df),
'严重等级分布': df['严重等级'].value_counts().to_dict(),
'质量维度分布': df['质量维度'].value_counts().to_dict(),
'涉及表分布': df['表名'].value_counts().to_dict()
}
缺陷收集器统一记录所有发现的问题,标准化数据结构,方便后续生成报告和统计分析。
3. 完整性校验
def check_completeness(df: pd.DataFrame, table_name: str,
required_cols: list, collector: DefectCollector):
for col in required_cols:
if col not in df.columns:
continue
# 同时检测 NaN 和空字符串
null_mask = df[col].isna() | (df[col].astype(str).str.strip() == '')
null_rows = df[null_mask]
for idx in null_rows.index:
collector.add(
rule_id='R01', dimension="完整性", table=table_name,
row_index=idx, column=col,
value=df.loc[idx, col] if pd.notna(df.loc[idx, col]) else '(空)',
description=f"{table_name}.{col} 为空,违反非空约束",
severity="高"
)
完整性校验同时处理 NaN 和空字符串,因为很多 CSV 文件中的空值可能表现为空字符串。
4. 一致性校验
def check_consistency(data: dict, collector: DefectCollector):
# R06: 订单的 user_id 必须在用户表中存在
valid_user_ids = set(data['users']['user_id'].dropna().unique())
for idx, row in data['orders'].iterrows():
uid = row['user_id']
if pd.notna(uid) and uid.strip() != '' and uid not in valid_user_ids:
collector.add(
rule_id='R06', dimension='一致性', table='orders',
row_index=idx, column='user_id', value=uid,
description=f"orders.user_id={uid} 在 users 表中不存在(孤儿订单)",
severity="高"
)
# R08: 订单状态值必须在合法枚举范围内
valid_statuses = {'pending', 'completed', 'shipped', 'cancelled', 'refunded'}
for idx, row in data['orders'].iterrows():
status = str(row['status']).strip().lower() if pd.notna(row['status']) else ''
if status and status not in valid_statuses:
collector.add(
rule_id='R08', dimension='一致性', table='orders',
row_index=idx, column='status', value=row['status'],
description=f"orders.status='{row['status']}' 不在合法集合中",
severity="中"
)
一致性校验包括参照完整性(外键关系)和枚举值一致性,是数仓数据质量保障的重点。
5. 准确性校验
def check_accuracy(data: dict, collector: DefectCollector):
# R10: 邮箱格式校验
email_pattern = r'^[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}$'
for idx, row in data['users'].iterrows():
email = str(row['email']).strip() if pd.notna(row['email']) else ''
if email and not re.match(email_pattern, email):
collector.add(
rule_id='R10', dimension='准确性', table='users',
row_index=idx, column='email', value=email,
description="邮箱格式不符合规范",
severity="中"
)
# R12: 订单金额必须 > 0
for idx, row in data['orders'].iterrows():
amount = float(row['amount']) if pd.notna(row['amount']) else None
if amount is not None and amount <= 0:
collector.add(
rule_id='R12', dimension='准确性', table='orders',
row_index=idx, column='amount', value=row['amount'],
description=f"orders.amount={amount},金额应大于0",
severity="高" if amount < 0 else "中"
)
6. 唯一性校验
def _check_duplicates(df: pd.DataFrame, column: str, table_name: str,
rule_id: str, collector: DefectCollector):
dup_mask = df.duplicated(subset=[column], keep=False)
dup_groups = df[dup_mask].groupby(column)
for dup_value, group in dup_groups:
for idx in group.index:
collector.add(
rule_id=rule_id, dimension='唯一性', table=table_name,
row_index=idx, column=column, value=dup_value,
description=f"{table_name}.{column}='{dup_value}' 重复出现,违反主键唯一约束",
severity="高"
)
df.duplicated(subset=[column], keep=False) 会标记所有重复的行,包括首次出现的那一行。
7. 报告生成
def generate_report(collector: DefectCollector, total_users: int,
total_orders: int, total_products: int) -> str:
df = collector.report()
summary = collector.summary()
total_records = total_users + total_orders + total_products
report = f"""
{'=' * 60}
电商数据质量测试报告
{'=' * 60}
生成时间:{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}
数据范围:users({total_users}行) | orders({total_orders}行) | products({total_products}行)
总缺陷数:{summary.get('总缺陷数', 0)}
总记录数:{total_records}
错误率: {summary.get('总缺陷数', 0) / total_records * 100:.2f}%
"""
return report
报告分为三部分:缺陷汇总、缺陷明细、修复建议,分别面向管理层、数据修复人员、架构师/开发。
8. 主流程
def main():
data_dir = Path(__file__).parent.parent / "data"
data = load_data(str(data_dir))
collector = DefectCollector()
# 按维度依次执行 19 条校验规则
check_completeness_all(data, collector) # R01 ~ R05
check_consistency(data, collector) # R06 ~ R09
check_accuracy(data, collector) # R10 ~ R14
check_timeliness(data, collector) # R15 ~ R16
check_uniqueness(data, collector) # R17 ~ R19
# 生成报告
report = generate_report(
collector,
total_users=len(data['users']),
total_orders=len(data['orders']),
total_products=len(data['products'])
)
# 保存 txt 报告和 Excel 缺陷明细
output_dir = Path(__file__).parent.parent / "output"
report_date = datetime.now().strftime('%Y%m%d_%H%M%S')
report_path = output_dir / f"data_quality_report_{report_date}.txt"
with open(report_path, 'w', encoding='utf-8') as f:
f.write(report)
defects_df = collector.report()
if len(defects_df) > 0:
excel_path = output_dir / f"defects_detail_{report_date}.xlsx"
defects_df.to_excel(excel_path, index=False, engine='openpyxl')
# 返回退出码,便于 CI/CD 集成
exit_code = 0 if collector.summary().get('总缺陷数', 0) == 0 else 1
return exit_code
典型脏数据案例
案例 1:订单金额为负数
| 字段 | 值 |
|---|---|
| order_id | ORD004 |
| amount | -50.00 |
| 违反规则 | R12 — 金额必须 > 0 |
影响:如果财务部门基于这个数据做月度结算,一笔 -50 元的订单会直接吃掉 5 笔 +10 元订单的利润,账就平不了。
案例 2:user_id 重复
| user_id | 姓名 |
|---|---|
| U009 | 王五 |
| U009 | 冯十一 |
| 违反规则 | R17 — 用户 ID 唯一性 |
影响:当用户表和订单表做 JOIN 时,一个 user_id 对应的订单会被错误地关联到另一个 user_id 的订单上,产生"张冠李戴"的数据污染,而且笛卡尔积会导致行数爆炸。
案例 3:孤儿订单
| order_id | user_id |
|---|---|
| ORD003 | U999(users 表中不存在) |
| 违反规则 | R06 — 参照完整性 |
影响:这部分订单无法归属到任何用户,用户画像分析、订单归因、客服查询全部失败。
项目成果
- 设计并实现 19 条数据质量校验规则,覆盖 5 大质量维度。
- 自动检测出 26 条数据缺陷(基于模拟脏数据),错误率 59.09%。
- 输出 txt 格式汇总报告和 Excel 格式缺陷明细。
- 提供按严重等级和质量维度的统计分析。
- 返回非 0 退出码,便于接入 CI/CD 流水线。
面试高频问题
Q1:数据质量测试测什么?
数据质量测试主要覆盖 DAMA 框架的五大维度:完整性、一致性、准确性、及时性、唯一性。完整性看字段是否缺失;一致性看跨表关系是否正确;准确性看数据格式和范围是否合理;及时性看时间字段是否正确;唯一性看主键是否重复。
Q2:为什么用 Pandas 做数据质量测试?
Pandas 是 Python 数据分析的瑞士军刀,可以方便地读取 CSV、JSON、Excel 等多种格式,支持向量化运算和灵活的过滤条件,非常适合做批量数据校验和统计分析。
Q3:发现脏数据后怎么处理?
一般分为三步:首先将脏数据拦截,不进入 DWD 层;然后写入异常数据表并触发告警;最后由数据团队人工修复,修复后重新跑校验。同时要从源头排查,比如前端校验、ETL 逻辑、数据库约束等。
写在最后
数据质量测试是数据工程和数据治理中非常重要的一环。本项目通过模拟真实电商场景,将理论知识与工程实践结合,帮助你建立系统化的数据质量测试思维。无论是面试数据测试、数据开发还是数据分析岗位,这些能力都会成为重要的加分项。