很多数据问题并不会让程序立刻报错。日期列里混入一个不存在的日期、年龄列里出现字符串、同一个业务主键重复两次,代码仍然可能继续运行,最后产出的报表却已经不可信。
比“发现错误后再清洗”更稳妥的做法,是在数据进入分析、建模或数据库之前设置一道质量闸门:每个批次先通过一组明确、可计算的规则,再决定正常放行、人工复核还是隔离。本篇用 pandas 完成一套可重复执行的行级检查,并把检查结果转换成图表与批次审计摘要。
解决的问题
输入是一张会员明细表,每行代表一名会员,包含会员编号、注册日期、最近下单日期、年龄、会员等级和近90天消费额。案例要同时检查四类常见问题:
最终输出不只是一句“数据有问题”,而是每条规则的失败数量、每行的失败原因、可直接进入下游的合格数据、等待处理的异常队列,以及一个可供自动任务读取的批次状态。本文使用固定随机种子生成模拟数据,所有数字只用于演示工作流,不代表真实业务分布。
方法原理
数据质量规则最适合写成与原表等长的布尔序列:True 表示该行违反规则,False 表示通过。把多条规则并排放进一个 DataFrame 后,按列求和得到每条规则的失败规模,按行求和则得到每条记录违反了多少条规则。
日期和数值转换是这套设计的关键。pd.to_datetime(..., errors="coerce") 会把无法解析的日期转成 NaT,pd.to_numeric(..., errors="coerce") 会把非法数值转成缺失值。这样,混杂在字段中的脏值不会让整个批次中断,而会被显式纳入失败规则。这里的“容错”不是静默忽略错误,而是先把错误标准化,再集中审计。
质量闸门还需要一个发布标准。本例把完全通过的记录标为 PASS,只违反一条规则的记录标为 REVIEW,同时违反两条及以上规则的记录标为 QUARANTINE。只有 PASS 记录进入下游;此外,批次完全通过率必须达到95%,否则批次状态为 BLOCK。这个阈值只是演示值,实际项目应根据字段重要性、历史基线和业务风险设定,不能机械照搬。
数据与实现思路
案例生成1200行模拟会员数据,再主动注入重复主键、非法日期、无法解析的年龄、越界年龄、未知等级、负消费额和时间顺序错误。检查流程分成四步:
哈希指纹不能证明数据内容正确,但能帮助判断输入或审计清单是否发生变化。同一份输入会得到相同的 SHA-256 摘要,因此它适合用于批次追踪和可重复性核对。
Python完整实战
1. 生成带有质量问题的模拟批次
import hashlibimport jsonimport randomfrom pathlib import Pathimport matplotlibmatplotlib.use(”Agg”)import matplotlib.pyplot as pltimport numpy as npimport pandas as pdimport seaborn as sns# 固定两套随机种子,让模拟数据和错误注入都可以复现random.seed(20260726)rng = np.random.default_rng(20260726)# 生成一批模拟会员数据,字段设计贴近常见业务明细表n = 1200base_date = pd.Timestamp(”2025-01-01”)signup_days = rng.integers(0, 540, n)signup = base_date + pd.to_timedelta(signup_days, unit=”D”)last_order = signup + pd.to_timedelta(rng.integers(0, 220, n), unit=”D”)# 金额采用右偏分布,避免把演示数据误写成近似正态raw = pd.DataFrame({ ”user_id”: [f”U{i:05d}” for i in range(1, n + 1)], ”signup_date”: signup.strftime(”%Y-%m-%d”), ”last_order_date”: last_order.strftime(”%Y-%m-%d”), ”age”: rng.integers(18, 76, n).astype(object), ”tier”: rng.choice([”Basic”, ”Silver”, ”Gold”], n, p=[0.56, 0.29, 0.15]), ”spend_90d”: np.round(rng.gamma(2.2, 420, n), 2),})# 注入六类数据问题;索引允许少量重叠,更符合真实批次raw.loc[rng.choice(n, 14, replace=False), ”user_id”] = raw.loc[:13, ”user_id”].to_numpy()raw.loc[rng.choice(n, 18, replace=False), ”signup_date”] = ”2026-02-30”raw.loc[rng.choice(n, 15, replace=False), ”age”] = ”unknown”raw.loc[rng.choice(n, 20, replace=False), ”age”] = rng.choice([12, 105], 20)raw.loc[rng.choice(n, 16, replace=False), ”tier”] = ”Diamond”raw.loc[rng.choice(n, 22, replace=False), ”spend_90d”] *= -1cross_idx = rng.choice(n, 19, replace=False)raw.loc[cross_idx, ”last_order_date”] = ”2024-12-01”print(f”rows={len(raw)}, columns={raw.shape[1]}”)print(f”declared_key=user_id, duplicate_rows={raw.duplicated('user_id', keep=False).sum()}”)print(”source=simulated customer batch, seed=20260726”)
输出:
rows=1200, columns=6declared_key=user_id, duplicate_rows=28source=simulated customer batch, seed=20260726
这里显示28个重复主键行,而不是注入时的14,原因是 duplicated(keep=False) 会同时标记重复组中的原始记录和后来产生的重复记录。这种写法适合生成异常队列,因为重复关系中的每一行都需要被核对。
2. 把数据要求转换为布尔规则
# 解析失败转成缺失值,使“脏值”变成可计算的布尔规则clean = raw.copy()clean[”signup_parsed”] = pd.to_datetime(clean[”signup_date”], errors=”coerce”)clean[”last_order_parsed”] = pd.to_datetime(clean[”last_order_date”], errors=”coerce”)clean[”age_parsed”] = pd.to_numeric(clean[”age”], errors=”coerce”)# 每一列代表一条可解释规则,True表示该行违反规则violations = pd.DataFrame(index=clean.index)violations[”Duplicate key”] = clean.duplicated(”user_id”, keep=False)violations[”Invalid signup date”] = clean[”signup_parsed”].isna()violations[”Invalid age type”] = clean[”age_parsed”].isna()violations[”Age out of range”] = clean[”age_parsed”].notna() & ~clean[”age_parsed”].between(18, 100)violations[”Unknown tier”] = ~clean[”tier”].isin([”Basic”, ”Silver”, ”Gold”])violations[”Negative spend”] = clean[”spend_90d”].lt(0)violations[”Order before signup”] = ( clean[”signup_parsed”].notna() & clean[”last_order_parsed”].notna() & clean[”last_order_parsed”].lt(clean[”signup_parsed”]))# 规则汇总既保留失败数,也给出便于批次比较的失败率rule_summary = ( violations.sum() .rename(”failed_rows”) .to_frame() .assign(failure_rate=lambda x: x[”failed_rows”] / len(clean)) .sort_values(”failed_rows”, ascending=False))print(rule_summary.assign( failure_rate=lambda x: x[”failure_rate”].map(”{:.2%}”.format)).to_string())# 第一张图直接呈现各规则失败规模,便于决定优先修复顺序sns.set_theme(style=”whitegrid”)fig, ax = plt.subplots(figsize=(9, 5.5))plot_data = rule_summary.sort_values(”failed_rows”)ax.barh(plot_data.index, plot_data[”failed_rows”], color=”#2F6690”)for i, value in enumerate(plot_data[”failed_rows”]): ax.text(value + 0.4, i, str(int(value)), va=”center”, fontsize=10)ax.set(title=”Data Quality Failures by Rule”, xlabel=”Failed rows”, ylabel=””)ax.set_xlim(0, plot_data[”failed_rows”].max() * 1.18)fig.tight_layout()fig.savefig(”/tmp/python_quality_rules.png”, dpi=180, bbox_inches=”tight”)plt.close(fig)
输出:
failed_rows failure_rateDuplicate key 28 2.33%Negative spend 22 1.83%Age out of range 20 1.67%Order before signup 19 1.58%Invalid signup date 18 1.50%Unknown tier 16 1.33%Invalid age type 14 1.17%
柱形图把规则从“代码条件”转换成了可排序的修复任务。本批次失败最多的是重复主键,其次是负消费额。失败数只能反映出现频率,不能直接代表业务损失;如果某个低频错误会改变付款、权限或监管结果,它仍然应当拥有更高优先级。
3. 识别问题共现并拆分数据去向
# 第二张图显示规则同时失败的次数,用于识别问题是否成簇出现cofailure = violations.astype(int).T.dot(violations.astype(int))mask = np.triu(np.ones_like(cofailure, dtype=bool), k=1)fig, ax = plt.subplots(figsize=(8.5, 7))sns.heatmap( cofailure, mask=mask, annot=True, fmt=”d”, cmap=”YlOrRd”, linewidths=0.5, cbar_kws={”label”: ”Rows”}, ax=ax,)ax.set_title(”Rule Co-failure Matrix”)ax.set_xlabel(””)ax.set_ylabel(””)fig.tight_layout()fig.savefig(”/tmp/python_quality_cofailure.png”, dpi=180, bbox_inches=”tight”)plt.close(fig)# 行级失败数把多条规则压缩成可用于分流的质量状态clean[”failed_rule_count”] = violations.sum(axis=1)clean[”quality_status”] = np.select( [clean[”failed_rule_count”].eq(0), clean[”failed_rule_count”].eq(1)], [”PASS”, ”REVIEW”], default=”QUARANTINE”,)# 只有完全通过的数据进入下游,其他记录保留原值和失败原因accepted = clean.loc[clean[”quality_status”].eq(”PASS”), raw.columns].copy()rejected = clean.loc[~clean[”quality_status”].eq(”PASS”), raw.columns].copy()rejected[”failed_rules”] = violations.loc[rejected.index].apply( lambda row: ”; ”.join(row.index[row]), axis=1)print(clean[”quality_status”].value_counts().reindex( [”PASS”, ”REVIEW”, ”QUARANTINE”], fill_value=0).to_string())print(f”accepted_rate={len(accepted) / len(clean):.2%}”)print(f”rows_with_multiple_failures={(clean['failed_rule_count'] >= 2).sum()}”)print( ”charts_created=” f”{Path('/tmp/python_quality_rules.png').stat().st_size > 0 and Path('/tmp/python_quality_cofailure.png').stat().st_size > 0}”)
输出:
quality_statusPASS 1067REVIEW 129QUARANTINE 4accepted_rate=88.92%rows_with_multiple_failures=4charts_created=True
热图对角线是每条规则自身的失败总数,对角线以下则是两条规则同时失败的行数。本次模拟中的大多数问题彼此独立,只有4行同时违反多条规则。若真实数据中某几个规则经常共同失败,通常意味着它们可能来自同一个录入界面、上游系统或转换步骤,修复源头比逐行修改更有效。
4. 生成机器可读的批次闸门
# 批次清单记录输入指纹、规则数量和是否放行,便于自动任务读取source_bytes = raw.to_csv(index=False).encode(”utf-8”)manifest = { ”batch_id”: ”customer_20260726”, ”input_rows”: int(len(raw)), ”accepted_rows”: int(len(accepted)), ”rejected_rows”: int(len(rejected)), ”rule_count”: int(violations.shape[1]), ”sha256”: hashlib.sha256(source_bytes).hexdigest(),}# 示例阈值要求至少95%的行完全通过;未达标时整批不自动发布manifest[”accepted_rate”] = round(len(accepted) / len(raw), 4)manifest[”gate_status”] = ”PASS” if manifest[”accepted_rate”] >= 0.95 else ”BLOCK”# JSON使用排序键生成稳定文本,同一输入可获得相同审计摘要manifest_text = json.dumps(manifest, ensure_ascii=False, sort_keys=True)manifest_digest = hashlib.sha256(manifest_text.encode(”utf-8”)).hexdigest()print(f”batch_id={manifest['batch_id']}”)print( f”gate_status={manifest['gate_status']}, ” f”threshold=95.00%, actual={manifest['accepted_rate']:.2%}”)print( f”accepted={manifest['accepted_rows']}, ” f”rejected={manifest['rejected_rows']}, rules={manifest['rule_count']}”)print(f”input_sha256={manifest['sha256'][:16]}...”)print(f”manifest_sha256={manifest_digest[:16]}...”)
输出:
batch_id=customer_20260726gate_status=BLOCK, threshold=95.00%, actual=88.92%accepted=1067, rejected=133, rules=7input_sha256=7b9e1e8b43211776...manifest_sha256=54d779d39ecb7d62...
闸门最终返回 BLOCK,因为完全通过率为88.92%,低于预设的95%。这不意味着必须丢弃整个文件:1067行合格记录已经被单独保存,133行异常记录也带有失败原因。BLOCK 的含义是禁止未经确认就把本批次整体发布到下游。
结果解释
这次运行最重要的结果不是发现了多少错误,而是把“发现—定位—分流—审计”连接成了一条完整路径。7条规则共标记133行异常,其中129行只违反一条规则,4行违反多条规则。完全通过率88.92%,因此批次未达到95%的自动放行标准。
规则失败数和异常行数不能简单相加。规则失败数统计的是“行—规则”组合,同一行可能同时贡献给多个规则;异常行数统计的是至少违反一条规则的独立记录。这个区别也是为什么规则矩阵和行级状态需要同时保留。
本例把所有规则视为同等重要,是为了讲清实现逻辑。生产环境通常还需要给规则分级。例如,展示名称为空可能只需提醒,付款金额为负或业务主键重复则可能直接阻断。此时可以加入 severity 字段,根据严重级别和失败数量共同决定 PASS、REVIEW 与 QUARANTINE。
另外,errors="coerce" 只是把解析失败显式转成缺失值,并没有自动修复原始数据。异常队列必须保留原始字段,修正过程也应记录负责人、修正时间和依据,否则“清洗后的正确值”仍然缺乏可追溯性。
结论
用 pandas 搭建数据质量闸门的核心,是把每条数据要求变成可解释的布尔规则,再同时做规则级汇总和行级分流。这样既能回答“哪类问题最多”,也能回答“哪些记录可以进入下游”。
这套方法适合每日文件导入、报表刷新、模型训练前检查和数据库入库前验证。它的边界也很明确:规则只能发现已经被定义的问题,无法代替业务知识,也无法证明通过检查的数据一定真实。高质量的数据流程需要持续根据事故、字段变化和上游系统调整规则与阈值。
当数据量和团队规模进一步增加时,可以把规则配置、批次清单和异常队列写入独立存储,并接入调度系统报警。不过无论使用什么框架,本文这套“解析、布尔规则、分流、审计”的基础结构仍然成立。