跳到内容

2.3 可复现清洗与质量门禁:每次改动都能解释和回滚

数据预言厅里,一份被手工修过的战报无法重现,档案管理员要求每条清洗规则都留下版本、证据和回滚路径。

手工改掉一份 CSV 很快,却无法回答三天后的问题:改了哪些行、用了哪版规则、原值是什么、重跑是否得到同样结果。清洗的专业边界不在 API 数量,而在可重现、可审计和可验证。

本课目标

  • 把清洗步骤设计成版本化转换;
  • 使用 quarantine 隔离坏记录而非静默丢弃;
  • 用计数、守恒量和属性测试建立质量门禁;
  • 分离 fit 参数与 transform 执行,防止泄漏。

1. 输入不可变,输出有版本

text
raw/source_version
  + code_version
  + rule_config_version
  + reference_data_version
  → curated/output_version

不要在原始表上原地修改。每次 run 记录输入清单、代码 commit/hash、配置、开始结束时间和输出校验。用户要求仓库不 commit 不影响教材里的通用版本概念;运行中的代码版本仍需可识别。

清洗函数尽量采用“输入 DataFrame,返回新 DataFrame + 质量报告”的接口,减少隐式全局状态。

2. 每条规则返回结果和原因

python
from dataclasses import dataclass
import pandas as pd

@dataclass(frozen=True)
class CleanResult:
    accepted: pd.DataFrame
    quarantined: pd.DataFrame
    metrics: dict[str, int | float]

def validate_temperature(frame: pd.DataFrame) -> CleanResult:
    invalid = ~frame["temperature_c"].between(-50, 150)
    quarantine = frame.loc[invalid].assign(
        rejection_code="TEMPERATURE_OUT_OF_RANGE"
    )
    accepted = frame.loc[~invalid].copy()
    return CleanResult(
        accepted=accepted,
        quarantined=quarantine,
        metrics={
            "input_rows": len(frame),
            "accepted_rows": len(accepted),
            "quarantined_rows": len(quarantine),
        },
    )

示例把非法温度整行隔离;真实系统也可能只把该字段置空。策略由下游需求决定,但不能静默删除。

隔离区需要同等级访问控制,因为其中往往含原始敏感数据。设置保留期、重处理状态和 owner,避免成为永久垃圾场。

3. Pipeline 步骤要有清晰顺序

一个可解释顺序:

  1. 解码与 schema 解析;
  2. 统一显式缺失标记;
  3. 逻辑类型转换;
  4. 单字段与跨字段有效性检查;
  5. 确定性去重;
  6. 标准化与衍生字段;
  7. 统计填补或模型变换;
  8. 输出契约与 reconciliation。

顺序会影响结果。若先把无效 999 温度用于计算中位数,再把它隔离,填补参数已经被污染。规则依赖应在管道定义中显式表达。

4. 行数核对是一条方程

若每条输入最终只能进入 accepted 或 quarantine:

$$ N_{input}=N_{accepted}+N_{quarantined}. $$

若还做去重:

$$ N_{input}=N_{accepted}+N_{quarantined}+N_{duplicates}, $$

前提是三个集合互斥。若 duplicate 也可能 invalid,需要规定分类优先级或使用多标签计数,不能硬套方程。

金额、事件数和实体数也可建立守恒检查。清洗改变数值时,输出原值与差异汇总。

5. 质量门禁比较绝对规则与基线

硬约束

  • 主键不得为空;
  • schema 必须兼容;
  • 目标分区不得缺失;
  • 引用键必须满足约定完整性。

统计监控

  • 行数相对季节基线;
  • 缺失率变化;
  • 类别新增与消失;
  • 分位数和分布漂移;
  • quarantine/retry 比例。

硬约束失败可阻断发布;统计偏移可能先告警、人工确认。把所有变化都阻断会导致频繁误报,把所有变化都放行又失去门禁价值。

6. fittransform 分开保存

python
from dataclasses import dataclass

@dataclass(frozen=True)
class TemperatureImputer:
    median: float

    @classmethod
    def fit(cls, training: pd.Series) -> "TemperatureImputer":
        return cls(median=float(training.median()))

    def transform(self, values: pd.Series) -> pd.Series:
        return values.fillna(self.median)

保存参数值、训练数据版本、特征定义和代码版本。线上若重新计算 median,会造成训练/服务偏差;回放历史也无法复现当时结果。

非机器学习分析同样需要版本化参数。例如异常阈值、汇率表和地区映射都会随时间变化。

7. 幂等性与确定性

理想清洗转换对同一输入和同一版本产生相同输出。检查:

  • 去重 tie-breaker 是否唯一;
  • 当前时间或随机数是否显式注入;
  • 外部维表是否固定版本;
  • 并行聚合的浮点顺序是否影响末位;
  • 输出排序是否属于契约。

有些转换还应满足幂等:

$$ clean(clean(x))=clean(x). $$

并非所有转换都天然幂等,例如每次运行都附加后缀或重复标准化。属性测试能迅速发现。

8. 测试从示例扩展到不变量

单元案例

覆盖合法、边界、缺失、未知枚举和解析失败。

属性测试

  • 输出 key 唯一;
  • 除明确策略外不增加行;
  • span/时间顺序不被颠倒;
  • 清洗两次结果相同;
  • accepted 与 quarantine 不重叠。

Golden dataset

保存小而有代表性的输入及期望输出,用于规则变更 diff。不要只断言文件 hash;输出变化时应展示具体行、列和原因。

Shadow run

新规则先并行运行但不发布,比较行数、指标和下游影响,再切换版本。

9. 发布报告说明改了什么

text
run_id: clean-2026-08-03-01
input: 1,000,000
accepted: 982,140
quarantine: 12,300
duplicates: 5,560
changed values:
  normalized location: 320,441
  imputed temperature: 8,201
top rejection reasons:
  invalid timestamp: 7,104
  impossible temperature: 5,196

比例突变应能下钻到来源、分区、schema 版本和规则。报告既服务运行人员,也为分析者解释数据版本之间的差异。

常见误区

  • 清洗脚本能重跑就算可复现:外部维表、参数和输入版本也必须固定。
  • 坏行直接删除最干净:会失去根因、审计和修复机会。
  • 测试几行示例足够:数据规则需要不变量、边界和真实分布回放。
  • 统计漂移必然是数据故障:也可能是真实业务变化,需要分级响应。

练习

  1. 把一个就地修改 DataFrame 的脚本改成 CleanResult 接口。
  2. 为缺失、异常和重复集合写 reconciliation 方程并处理交集。
  3. 给清洗函数写 clean(clean(x)) == clean(x) 属性测试。
  4. 设计一份新旧规则 shadow run 差异报告。

小结

可靠清洗是版本化的数据转换:原始输入不变,规则命中有原因,坏记录可隔离,输出满足不变量,参数只能从允许的数据拟合。这样清洗才从一次性 notebook 操作变成可审计的生产环节。

下一章使用清洗后的数据做 EDA。第一步仍不是画漂亮图,而是明确每张图在检查分布、关系还是数据生成过程。

Built with VitePress | Software Systems Atlas