数据测试自动化复盘:从手工核对到自动校验的工程化过程
数据测试自动化复盘从手工核对到自动校验的工程化过程一、手工核对阶段的痛苦回忆数据测试听起来不像个严肃话题——不就是看看数据对不对嘛但做过数据项目的人都知道数据测试的痛苦程度远超代码测试。代码有bug会抛异常数据有问题不会抛任何东西它只是悄悄地错了然后下游报表和模型都基于错误数据运行得出错误的结论。我们团队的手工核对阶段持续了将近一年。那时候的数据测试流程是这样的ETL任务跑完后分析师手动打开数据库执行几条SELECT语句看看数据量是否正常关键报表数据出来后手动跟上周数据对比看波动是否合理新上线的数据管道手动造一批测试数据跑一遍看看结果是否符合预期这个流程的问题很明显看这个流程就知道了——肉眼判断是否正常和下游才发现数据错误是两个最大的漏洞。肉眼判断疲劳后出错率飙升下游发现意味着错误已经扩散了。手工核对的典型失误案例7月份有个案例一个新上线的用户画像ETL任务把user_age字段从年龄段改成了具体年龄值但下游的消费预测模型还在用年龄段做特征——模型跑了一个星期才被发现输入特征变了预测结果偏差严重但没人注意到因为偏差是渐进式的不是突然的。二、数据测试的分类体系从手工核对转向自动化之前我们先定义了数据测试的分类体系。数据测试不像代码测试有清晰的单元测试/集成测试分层它需要按数据生命周期来分类各类测试的具体规则定义from enum import Enum from dataclasses import dataclass from typing import List, Callable class TestCategory(Enum): 数据测试类别枚举 SOURCE_DATA 源数据测试 ETL_PROCESS ETL过程测试 OUTPUT_DATA 产出数据测试 REPORT_MODEL 报表/模型测试 class Severity(Enum): 测试严重级别 CRITICAL critical # 阻断性测试失败则ETL任务中断 HIGH high # 高测试失败发告警任务继续 MEDIUM medium # 中测试失败仅记录 LOW low # 低测试失败仅统计 dataclass class DataTestRule: 数据测试规则定义 test_id: str # 测试规则ID name: str # 测试名称 category: TestCategory # 测试类别 severity: Severity # 严重级别 target_table: str # 目标表名 check_fn: Callable # 校验函数 threshold: float # 阈值如非空率不低于0.95 description: str # 测试描述三、自动化测试框架的搭建有了分类体系开始搭建自动化测试框架。我们的框架设计参考了 great_expectations 的理念但做了大量简化和定制。核心测试函数库import pandas as pd import numpy as np class DataTestEngine: 数据测试自动化引擎 def __init__(self): self.test_results [] def check_not_null_ratio(self, df, column, threshold0.95): 检查字段非空率 参数: df: 数据DataFrame column: 待检查的列名 threshold: 非空率最低阈值 actual_ratio 1 - df[column].isna().mean() passed actual_ratio threshold self.test_results.append({ test: f{column} 非空率 {threshold}, actual: f{actual_ratio:.4f}, passed: passed }) return passed def check_row_count_range(self, df, min_rows, max_rows): 检查数据行数是否在合理范围内 参数: df: 数据DataFrame min_rows: 最小行数 max_rows: 最大行数 actual_rows len(df) passed min_rows actual_rows max_rows self.test_results.append({ test: f行数在 [{min_rows}, {max_rows}] 范围内, actual: str(actual_rows), passed: passed }) return passed def check_value_range(self, df, column, min_val, max_val): 检查字段值是否在合理范围内 参数: df: 数据DataFrame column: 待检查的列名 min_val: 最小值 max_val: 最大值 col_min df[column].min() col_max df[column].max() passed col_min min_val and col_max max_val self.test_results.append({ test: f{column} 值范围在 [{min_val}, {max_val}], actual: f实际范围 [{col_min}, {col_max}], passed: passed }) return passed def check_unique_ratio(self, df, column, threshold1.0): 检查字段唯一值比例主键字段应为1.0 参数: df: 数据DataFrame column: 待检查的列名 threshold: 唯一值比例阈值 actual_ratio df[column].nunique() / len(df) passed actual_ratio threshold self.test_results.append({ test: f{column} 唯一率 {threshold}, actual: f{actual_ratio:.4f}, passed: passed }) return passed def check_referential_integrity(self, df, column, ref_df, ref_column): 检查外键引用完整性子表中的值必须存在于父表中 参数: df: 子表DataFrame column: 子表外键列 ref_df: 父表DataFrame ref_column: 父表主键列 child_values set(df[column].unique()) parent_values set(ref_df[ref_column].unique()) orphan_count len(child_values - parent_values) passed orphan_count 0 self.test_results.append({ test: f{column} 外键完整性, actual: f孤立记录数: {orphan_count}, passed: passed }) return passed def check_statistical_change(self, df, column, hist_df, hist_column, max_pct_change0.3): 检查统计指标的波动是否在合理范围内 参数: df: 当前数据 column: 当前数据列名 hist_df: 历史参考数据 hist_column: 历史数据列名 max_pct_change: 最大允许变化比例 current_mean df[column].mean() hist_mean hist_df[hist_column].mean() pct_change abs(current_mean - hist_mean) / hist_mean passed pct_change max_pct_change self.test_results.append({ test: f{column} 均值变化不超过 {max_pct_change:.0%}, actual: f变化 {pct_change:.2%} (当前 {current_mean:.2f} vs 历史 {hist_mean:.2f}), passed: passed }) return passed测试流水线的编排def run_data_test_pipeline(table_name, date): 运行数据测试流水线针对某张表执行全部相关测试 参数: table_name: 待测试的表名 date: 数据日期 engine DataTestEngine() # 加载待测试数据 df pd.read_parquet(f/data/output/{table_name}_{date}.parquet) # 加载历史参考数据前7天均值 hist_df pd.read_parquet(f/data/output/{table_name}_hist7d.parquet) # 源数据测试 # 核心字段非空率不低于99% engine.check_not_null_ratio(df, user_id, threshold0.999) engine.check_not_null_ratio(df, order_id, threshold0.999) engine.check_not_null_ratio(df, order_amount, threshold0.99) # 数据量在合理范围 engine.check_row_count_range(df, min_rows50000, max_rows500000) # ETL过程测试 # 主键唯一性 engine.check_unique_ratio(df, order_id, threshold1.0) # 外键完整性订单的用户ID必须存在于用户表中 user_df pd.read_parquet(/data/output/user_table.parquet) engine.check_referential_integrity(df, user_id, user_df, user_id) # 产出数据测试 # 订单金额在合理范围内0-100000 engine.check_value_range(df, order_amount, min_val0, max_val100000) # 核心指标波动不超过30% engine.check_statistical_change(df, order_amount, hist_df, order_amount, max_pct_change0.3) # 报表/模型测试 # 检查特征字段类型是否一致防止偷偷从年龄段改成具体年龄值 if user_age in df.columns: age_dtype df[user_age].dtype expected_dtype int64 # 期望是整数年龄段 engine.test_results.append({ test: user_age 字段类型为 int64, actual: str(age_dtype), passed: age_dtype expected_dtype }) # 输出测试报告 report pd.DataFrame(engine.test_results) pass_rate report[passed].mean() print(f\n{*50}) print(f数据测试报告: {table_name} ({date})) print(f总测试数: {len(report)}, 通过: {report[passed].sum()}, 失败: {(~report[passed]).sum()}) print(f通过率: {pass_rate:.1%}) print(f{*50}) # 输出失败测试详情 failed_tests report[~report[passed]] if len(failed_tests) 0: print(\n失败测试详情:) for _, row in failed_tests.iterrows(): print(f ✗ {row[test]}: 实际值 {row[actual]}) return report四、从单表测试到全链路校验单表测试只是起点。数据管道是一个链路——源表→中间表→产出表→报表/模型一个环节的错误会逐级传导。我们需要全链路校验确保错误在上游就能被拦截。全链路校验的设计全链路校验的关键原则阻断性测试失败就停非阻断性测试失败继续但告警。from enum import Enum class TestAction(Enum): 测试失败时的处置动作 BLOCK 阻断 # 停止下游任务 ALERT 告警 # 发送告警但继续运行 LOG 仅记录 # 只记录不告警 def run_full_pipeline_validation(pipeline_config, date): 全链路校验按管道配置依次执行各环节测试 参数: pipeline_config: 管道配置包含各环节的测试规则列表 date: 数据日期 pipeline_status PASS for stage in pipeline_config[stages]: stage_name stage[name] print(f\n--- 执行 {stage_name} 测试 ---) # 执行该环节的所有测试 report run_data_test_pipeline(stage[table], date) # 检查是否有阻断性测试失败 for test_rule in stage[test_rules]: if test_rule.severity Severity.CRITICAL: # 在报告中查找对应的测试结果 test_result report[report[test] test_rule.name] if len(test_result) 0 and not test_result[passed].values[0]: print(f✗ 阻断性测试失败: {test_rule.name}) print(f → 停止下游任务执行) pipeline_status BLOCKED return pipeline_status, report # 检查是否有高严重级别测试失败 for test_rule in stage[test_rules]: if test_rule.severity Severity.HIGH: test_result report[report[test] test_rule.name] if len(test_result) 0 and not test_result[passed].values[0]: send_alert(test_rule, test_result) pipeline_status ALERTED print(f\n全链路校验结果: {pipeline_status}) return pipeline_status, report测试覆盖率统计我们从手工核对转向自动化后测试覆盖率大幅提升维度手工核对阶段自动化阶段测试规则数15条87条覆盖表数5张核心表23张表执行频率每日人工1次每次ETL自动执行平均耗时45分钟3分钟错误发现时间下游发现平均2天上游拦截实时87条测试规则覆盖23张表执行时间从45分钟降到3分钟——这不是人少了是机器把重复劳动接管了。五、总结数据测试从手工核对到自动校验的工程化过程复盘核心经验手工核对的致命缺陷是靠人判断人疲劳后判断力下降渐进式变化肉眼察觉不到下游发现意味着错误已经扩散测试分类体系是框架的基础按数据生命周期分四类源数据/ETL过程/产出数据/报表模型每类有不同的测试重点和严重级别阻断性测试 vs 告警性测试要区分清楚关键字段的非空率、主键唯一性、外键完整性这些是阻断性的必须失败就停统计波动类的是告警性的失败继续但发通知全链路校验比单表测试价值大10倍上游拦截一个错误比下游发现一个错误省掉整个排查和修复链路的时间字段类型检测是被低估的测试规则那个user_age从年龄段变成具体年龄值的案例就是字段类型偷偷变了导致的。类型检测规则虽然简单但效果惊人数据测试自动化的终点不是100%覆盖所有规则而是关键路径上的阻断性测试100%覆盖其余路径告警性测试逐步完善。先拦住最致命的错误再慢慢铺开覆盖率。