数据质量监控项目复盘:从 0 到 1 搭建质量保障体系

发布时间:2026/7/22 11:05:19
数据质量监控项目复盘:从 0 到 1 搭建质量保障体系 数据质量监控项目复盘从 0 到 1 搭建质量保障体系大家好我是朱大喜今天来复盘一个我花了大半年推进的项目——数据质量监控体系。记得项目刚启动的时候运营同学对我说大喜这周你们给的活动数据好像不太对啊……那一刻我就知道数据质量问题不能再靠用户反馈驱动了。一、问题背景数据质量问题到底有多严重在搭建质量监控体系之前我们先做了一次数据质量盘点。不盘不知道一盘吓一跳。我们统计了过去三个月的数据事故发现平均每月有 12 起由上下游感知到的数据异常问题。问题的类型五花八门上游 ET 任务延迟导致报表数据为空、数据倾斜导致某些分区的数据量异常、业务方改了枚举值没通知导致下游 JOIN 全丢、日志采集丢包导致部分用户行为缺失……最严重的一次运营同学用错误的活动数据做了一周的运营决策直到用户投诉才发现数据不对。每次出问题复盘会上的结论几乎都是我们能不能在数据出问题的时候自动发现而不是等业务方来投诉于是数据质量监控项目正式立项。在需求分析阶段我们梳理了六类最核心的质量问题完整性数据有没有丢表分区是否齐全准确性数据值对不对字段是否在合理范围内一致性同一份数据在不同地方是否一致及时性数据是否按时产出唯一性主键是否重复有效性数据格式是否符合规范二、系统架构设计考虑到团队的技术栈和现有基础设施我们选择了轻量规则引擎 Python 调度 企业微信告警的架构。核心思路是先保证能用、好维护不做大而全的质量平台而是做成一个质量插件嵌入到现有的 Airflow 调度链路中。整体架构分为四层检测层负责执行各种质量规则调度层由 Airflow 编排在每条 ETL 链路的关键节点插入质量检测任务告警层负责将问题推送给对应 Owner存储层记录每一次检测结果用于质量趋势分析。三、核心实现质量规则执行引擎质量规则引擎是整个系统的核心。我们设计了一个配置驱动的模式在 MySQL 中维护一张质量规则配置表每个规则对应一条记录。Python 脚本读取规则配置动态生成 SQL 执行检测。import pymysql import json from datetime import datetime import requests # 质量规则配置读取与执行 class DataQualityChecker: 数据质量规则执行引擎 def __init__(self, db_config): self.db_config db_config self.conn None def connect(self): 建立数据库连接 self.conn pymysql.connect(**self.db_config) def load_rules(self, table_name, check_typeNone): 从规则配置表加载质量规则 支持按表名和检测类型过滤 sql SELECT rule_id, rule_name, table_name, check_type, check_sql, expect_value, compare_operator, severity, owner, is_enabled FROM data_quality_rules WHERE table_name %s AND is_enabled 1 params [table_name] if check_type: sql AND check_type %s params.append(check_type) with self.conn.cursor(pymysql.cursors.DictCursor) as cursor: cursor.execute(sql, params) return cursor.fetchall() def execute_check(self, rule): 执行单条质量规则 rule_id rule[rule_id] rule_name rule[rule_name] check_sql rule[check_sql] expect_value rule[expect_value] compare_operator rule[compare_operator] severity rule[severity] try: # 执行检测 SQL返回实际值 with self.conn.cursor() as cursor: cursor.execute(check_sql) result cursor.fetchone() actual_value result[0] if result else None # 执行比较逻辑 passed False if compare_operator and actual_value is not None: passed actual_value expect_value elif compare_operator and actual_value is not None: passed actual_value expect_value elif compare_operator and actual_value is not None: passed actual_value expect_value elif compare_operator ! and actual_value is not None: passed actual_value ! expect_value # 记录检测结果 status PASS if passed else FAIL # 失败时发送告警 if not passed: self.send_alert(rule_name, actual_value, expect_value, severity) # 写入质量检测日志 self.log_result(rule_id, actual_value, expect_value, status) return {rule_id: rule_id, rule_name: rule_name, status: status, actual: actual_value} except Exception as e: # 检测 SQL 执行异常也要记录和告警 error_msg f规则 [{rule_name}] 执行异常: {str(e)} print(error_msg) self.log_result(rule_id, None, expect_value, ERROR) return {rule_id: rule_id, rule_name: rule_name, status: ERROR, error: str(e)} def send_alert(self, rule_name, actual, expected, severity): 企业微信告警通知 根据严重程度分级通知 if severity CRITICAL: mentioned_list [all] # 严重问题 所有人 else: mentioned_list [] # 普通问题只发群消息 alert_msg f【数据质量告警 - {severity}】 检测规则{rule_name} 期望值{expected} 实际值{actual} 检测时间{datetime.now().strftime(%Y-%m-%d %H:%M:%S)} 请相关同学尽快排查 # 企业微信 Webhook 推送 webhook_url https://qyapi.weixin.qq.com/cgi-bin/webhook/send?keyYOUR_KEY payload { msgtype: text, text: { content: alert_msg, mentioned_list: mentioned_list } } requests.post(webhook_url, jsonpayload) def log_result(self, rule_id, actual_value, expect_value, status): 记录检测结果到日志表用于趋势分析 sql INSERT INTO data_quality_check_log (rule_id, actual_value, expect_value, status, check_time) VALUES (%s, %s, %s, %s, NOW()) with self.conn.cursor() as cursor: cursor.execute(sql, (rule_id, str(actual_value), str(expect_value), status)) self.conn.commit() def run_table_checks(self, table_name): 对指定表执行所有启用的质量规则 rules self.load_rules(table_name) results [] for rule in rules: result self.execute_check(rule) results.append(result) # 汇总输出 pass_count sum(1 for r in results if r[status] PASS) fail_count sum(1 for r in results if r[status] FAIL) error_count sum(1 for r in results if r[status] ERROR) print(f\n[{table_name}] 质量检测完成通过 {pass_count}, f失败 {fail_count}, 异常 {error_count}) return results # 使用示例 if __name__ __main__: checker DataQualityChecker({host: ***, user: ***, password: ***, database: data_quality}) checker.connect() # 对某个核心表执行全部质量检测 results checker.run_table_checks(dwd_order_detail) # 输出失败项详情 failed [r for r in results if r[status] in (FAIL, ERROR)] if failed: print(f\n以下 {len(failed)} 条规则未通过) for item in failed: print(f - {item[rule_name]}: {item[status]})四、规则配置与落地实践规则配置是这套系统的灵魂。我们从最容易实现的规则开始逐步覆盖。首批上线的规则很简单分区行数校验。用历史 7 天的行数均值和标准差判断当天产出的数据量是否在合理范围内。比如某张表日均 50 万行标准差 5000 行那当天行数如果在 48 万到 52 万之间就认为正常。第二波加上了主键唯一性和字段非空率检测。主键重复是数据加工中最常见的 Bug比如 JOIN 时笛卡尔积导致数据量翻倍、去重逻辑写错了导致重复插入。主键唯一性检测能第一时间兜住这类问题。第三波是分布漂移检测。用 KL 散度或 JS 散度来对比当天数据分布和历史分布的差异。这个比较高级主要用在一些核心业务指标源表上比如订单金额分布如果忽然向左偏小单变多可能是黑产薅羊毛。落地过程中最头疼的问题其实是误报处理。质量规则的阈值如果设得太紧天天告警满天飞团队很快就麻木了设得太松又怕漏报。我们的经验是先设一个宽松阈值跑一周观察实际波动范围然后逐步收紧。同时引入告警聚合机制——同一个表连续 3 次不通过才发送告警大幅降低了噪点。五、总结从 0 到 1 搭建数据质量监控体系技术实现的难度其实不大难的是如何在团队中建立数据质量文化。一开始大家会觉得天天告警烦死了但随着误报率降低、真问题被及时发现业务方和开发同学都慢慢认可了这套体系的价值。现在我们的质量监控覆盖了 200 张核心表、600 条规则月均自动发现 35 个数据问题其中 70% 在业务方感知前就被修复了。最重要的是再也没有运营同学跑来跟我说数据好像不太对了——因为他们已经收到了系统自动推送的数据健康日报。下次连载我们就来聊聊这个数据健康日报是怎么用 AI 自动生成的