做数仓的兄弟应该都经历过这种深夜打来的电话:报表跑完了,但一看数据明显不对,订单量少了三分之一,查了半天定位到上游同步任务凌晨挂了,重跑之后才发现脏数据已经污染了下游。数据量越大、链路越长,这种问题就越难靠人肉巡检兜住。所以这两年我一直在推进一件事:用开源工具拼一套数据质量监控平台,让规则自动跑、结果自动存、异常自动喊人。
这篇文章就把这套方案的完整搭建过程写出来。文章不侧重某个商业产品,而是围绕“一站式”这个目标,把开源生态里几个能打的数据质量工具按场景组合起来:离线大表用Deequ做深度校验,业务表快速接入用Soda Core做轻量检查,跨源数据对账用Apache Griffin,调度统一走DolphinScheduler,结果统一入库,告警统一推送,展示统一做Grafana看板。适合数仓工程师、数据平台开发、数据治理方向的同学参考,也适合团队刚准备做数据质量、还没想清楚怎么落地的朋友拿去直接用。
1. 先从业务痛点说起:为什么要自建数据质量平台
1.1 数据质量问题的真实代价
很多人一开始觉得“数据质量”是个很虚的概念,直到线上出了事故才意识到它的分量。最常见的场景是ETL链路里某个字段映射错了,或者源系统数据库扩容导致同步任务凌晨挂了然后自动重跑,重跑出来的结果跟正常业务时间对不上。这类问题往往不是立刻引爆的,而是等下游报表跑完、业务方早上拿到数据时才发现异常,这时候已经过去了两三个小时。
数据质量问题带来的直接代价不只是“数据不准”,还有信任崩塌。业务方会因为一次错误数据对整个数仓产生怀疑,之后每一张报表都要反复确认数据来源。这种信任一旦丢失,数据平台的生存空间会被大大压缩。更麻烦的是,当数据量从千万级涨到亿级、表数量涨到几千张时,纯靠人工去抽数、查数、盯数根本不可能,必须有系统化的工具在无人值守状态下持续运行。
1.2 平台需要解决的核心问题
我整理了一下,一个能真正落地的数据质量监控平台至少要解决四件事:
- 规则可配置:不同表、不同字段的质量标准不一样,订单表要求订单号不能为空,用户表要求邮箱格式正确,日活表要求行数波动不超过10%。规则必须能按表、按字段灵活配置,最好能用声明式写法,不用开发写一堆Java代码。
- 调度可管理:每天凌晨离线任务跑完后自动触发质量校验,校验任务要能被统一调度、重跑、补数,还能看到每个任务跑了多久、是否成功。
- 结果可追溯:某一天的数据质量不合格,要能查清楚是哪张表、哪个字段、哪个规则出了问题,当时的校验结果是什么样的,不能等业务方投诉了才去翻日志。
- 告警可送达:规则失败后要第一时间通知到责任人,而且不能一告警就刷屏,要有分级和分组策略。
如果用商业产品,比如Informatica、Collibra,功能确实全,但价格不是一般团队能承受的。用开源工具自己拼,成本低、可控性强,缺点是得自己踩坑。我下面要讲的方案,就是围绕这四件事去搭建的。
1.3 为什么选择开源工具而不是自研
不少团队一开始想自研数据质量平台,理由是“我们数据量特殊,只有自己写的才符合需求”。这个思路我不反对,但强烈建议先站在开源工具的肩膀上起步。原因很简单:数据质量的核心在于规则引擎和指标计算,这块的通用性远比你想象的高。完整性、唯一性、及时性、准确性,不管什么行业都逃不开这几个维度,开源工具已经沉淀了很多经验。
自研需要投入研发人力做规则解析、做任务调度、做结果存储、做告警通知,这些都属于“轮子”,而且做出来的稳定性大概率不如经过社区考验的成熟工具。我建议的路线是:先用开源工具把核心链路跑通,如果后续有非常个性化的需求,再在开源工具之上做扩展开发,而不是从头写一套。
2. 平台整体架构与开源工具选型
2.1 常见开源数据质量工具对比
目前社区里主流的开源数据质量工具主要有四个,我把它们放在一起对比过,各自的侧重点很不一样。
| 工具 | 底层引擎 | 擅长场景 | 规则定义方式 | 社区活跃度 |
|---|---|---|---|---|
| Apache Griffin | Spark | 跨源数据对账、批流一体质量监控 | JSON配置 | 中低 |
| Deequ | Spark | 离线大表多维度校验、数据画像 | Scala/Python API | 中高 |
| Soda Core | Python | 业务库/数仓表快速接入、轻量级检查 | YAML声明式 | 高 |
| Great Expectations | Python | 数据管道测试、数据文档 | Python DSL | 高 |
如果只看名字容易懵,我用一个类比来解释:Deequ更像一个资深数据科学家,适合对大规模数据集做非常细致的“体检”,能输出很多统计指标;Soda Core更像一个运维巡检员,配置简单、上手快,适合日常值班盯表;Griffin更像一个审计专员,专门做两边数据对不上号的问题;Great Expectations则更适合在数据管道里做单元测试。
2.2 我的选型结论:混合方案而不是押注单一工具
实话实说,没有任何一个开源工具能包打天下。Griffin适合做“源表和目标表数据量是否一致、对账是否平”这种跨源场景,但项目活跃度这些年明显下降,新版本迭代慢。Deequ功能强,可它是基于Spark的,每次跑批要起Spark任务,对平台资源有一定要求,不适合做高频率的简单检查。Soda Core很轻,安装简单,YAML写规则也很清晰,但复杂的数据血缘校验、大规模全量数据扫描它做不了。
所以我最后的方案是混合架构:Deequ做深度体检,Soda Core做日常巡检,Griffin做专项对账,DolphinScheduler做统一调度,Grafana做统一展示。这样既避免了单一工具的短板,又能让不同场景用最合适的工具。所谓“一站式”,不是只装一个软件,而是把工具链打通:规则一处配置、任务统一调度、结果集中存储、告警统一出口。
2.3 整体架构的六层拆解
整个平台架构我习惯分成六层来看,后续所有部署和操作都围绕这六层展开。
- 数据源层:包括业务数据库(MySQL、PostgreSQL)、数仓表(Hive、Iceberg、Hudi)、消息队列实时数据(Kafka)等,它们是质量校验的对象。
- 校验引擎层:由Soda Core、Deequ、Griffin组成,负责执行各种质量规则,输出校验结果。
- 调度层:使用DolphinScheduler统一调度校验任务,设定执行时间、重试策略和依赖关系。
- 存储层:校验结果统一写入PostgreSQL,作为质量元数据中心,同时存储规则配置和告警记录。
- 服务层:负责结果解析、异常判断、告警触发,一般是几个Python脚本或API服务。
- 展示层:Grafana连接质量结果库,按表、按日期、按规则维度展示通过率和波动趋势。
这套架构的好处是每一层都可以独立替换。比如调度层你不想用DolphinScheduler,换成Airflow或者简单的crontab也能跑;展示层不想用Grafana,用Superset也可以。核心思想是解耦。
3. 环境准备与核心组件安装部署
3.1 环境规划与前置条件
我先说一下我自己搭建时的环境,方便你对照参考。我这边是一个中型数仓集群,规模大概是20台服务器,Hadoop 3.3.4、Spark 3.2.1、Hive 3.1.3,调度用DolphinScheduler 3.1.8,质量结果库用PostgreSQL 13,装在独立的一台机器上。数据质量监控平台的组件单独部署在3台节点上,跟生产集群物理隔离,避免影响正常业务。
如果你们的集群规模比较小,完全可以把所有组件部署在同一台机器上,开发测试毫无问题。Soda Core是纯Python服务,安装很轻;Deequ是Spark库,只要提交到集群即可;DolphinScheduler只是调度,可以部署在边缘节点。
3.2 Soda Core部署:最简单,先跑通
Soda Core的官方名称为SodaCL(Soda Checks Language),核心是用YAML写检查规则。安装非常简单,使用pip安装即可。如果你的数据源是PostgreSQL,就安装对应的插件包:
python3 -m venv /opt/soda/venv source /opt/soda/venv/bin/activate pip install soda-core-postgres如果是MySQL或者其他数据源,把postgres换成对应的名称即可。装完验证一下版本:
soda --version接下来在项目目录下写两个文件:一个是数据源连接配置configuration.yml,一个是检查规则文件checks.yml。先看configuration.yml:
data_source order_dw: type: postgres host: 192.168.1.20 port: 5432 username: quality_ro password: your_password database: order_warehouse schema: public这里我建议给质量监控创建只读账号,永远不要用管理员账号连数据源,避免误操作。
3.3 DolphinScheduler部署:调度层接入
如果你团队里已经有DolphinScheduler或Airflow,直接复用就行。没有的话我建议先用最简单的方式起步,因为质量校验任务的执行逻辑往往很简单,就是定时跑一个命令或脚本。我在初期甚至直接用了crontab,每半小时跑一次核心表的检查,跑通之后才把任务迁到DolphinScheduler统一管理。
DolphinScheduler部署起来也不复杂,官方提供docker-compose方式可以快速起一套环境。生产环境我推荐用二进制包部署,在Apache DolphinScheduler官网下载对应版本,解压后修改conf/common.properties里的数据库连接,然后执行脚本初始化元数据库,最后分别启动master和worker服务。核心配置项如下:
# conf/common.properties spring.datasource.url=jdbc:postgresql://192.168.1.10:5432/dolphinscheduler spring.datasource.username=ds_user spring.datasource.password=ds_password注意DolphinScheduler的元数据库默认支持MySQL和PostgreSQL,建议独立部署,不要跟质量结果库混用。
3.4 质量结果表初始化
不管校验结果来自Soda Core还是Deequ,最终都会落在PostgreSQL里。我设计了一张统一的结果表,字段如下:
CREATE TABLE data_quality_check_result ( id BIGSERIAL PRIMARY KEY, exec_date DATE NOT NULL, table_name VARCHAR(128) NOT NULL, check_name VARCHAR(256) NOT NULL, check_status VARCHAR(32) NOT NULL, metric_name VARCHAR(64), metric_value NUMERIC, threshold VARCHAR(128), rule_version VARCHAR(32), batch_id VARCHAR(64), create_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX idx_quality_result_date_table ON data_quality_check_result(exec_date, table_name);这张表设计的时候我特意加了batch_id字段,用来关联某一次调度任务的执行。有了它,重跑任务时不会造成结果混淆,也能追踪到某次异常到底是由哪个批次的校验产生的。
4. 规则配置实操:以订单明细表为例
4.1 定义质量校验目标
好的规则定义是整个平台成败的关键。规则不是越多越好,而是要把关键表的“关键特征”抓准。我以数仓里最常见的ods_order_detail订单明细表为例,说一下我会从哪些维度定义规则。
这张表每天一个分区,每天凌晨3点由业务库同步到数仓。业务上最担心的问题包括:同步任务中途失败导致数据缺失、业务库表结构变更导致字段映射错误、重复同步导致数据翻倍。对应到数据质量维度,我至少要监控以下五项:
- 完整性:关键字段
order_id、user_id、order_amount不能为空,空值率必须为0。 - 唯一性:
order_id不能重复,重复订单会导致下游汇总翻倍。 - 及时性:当天分区数据的时间戳应与当前时间接近,如果数据延迟超过12小时,说明同步链路有问题。
- 准确性:
order_amount金额必须大于0,订单状态字段必须在规定的枚举值范围内。 - 稳定性:当日分区行数与前一天的波动幅度不能超过20%,超过则大概率是同步异常。
4.2 用Soda Core写巡检规则
Soda Core的YAML规则放到checks.yml里,写法非常直白。实际配置如下:
checks for ods_order_detail: - missing_count(order_id) = 0: name: 订单号完整性检查 severity: error - duplicate_count(order_id) = 0: name: 订单号唯一性检查 severity: error - missing_percent(user_id) < 1: name: 用户ID缺失率检查 severity: warning - freshness(order_time) < 24h: name: 数据新鲜度检查 - row_count > 0: name: 空表检查字段severity有两个级别:error和warning。我个人建议核心表全部用error级别,因为校验失败就要立刻阻断下游任务;普通维度表用warning级别,只记录不阻断,避免告警轰炸。freshness是Soda Core内置的“新鲜度”检查,它计算表中最大时间字段与当前时间的差值,如果超过24小时就认为数据不够新。
执行检查的命令如下:
cd /opt/quality_monitor /opt/soda/venv/bin/soda scan -d order_dw -c configuration.yml checks.yml --output-dir results/$(date +%F)执行完成后,results目录下会生成JSON格式的检查结果,包含每个检查的状态、指标值和错误信息。这个JSON就是后续入库和告警的数据源。
4.3 用Deequ做深度体检
Soda Core适合快速巡检,但要对全量数据做更深入的统计分析,比如字段值分布、指定列唯一值比例、跨表关联完整性,就需要Deequ上场了。Deequ是AWS开源的基于Spark的库,通过编写Scala或Python代码来定义约束,提交到Spark集群执行。
我用Python写的Deequ检查脚本片段如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import col import deequ from deequ import VerificationSuite, CheckLevel, Check spark = SparkSession.builder \ .appName("DeequQualityCheck") \ .enableHiveSupport() \ .getOrCreate() data = spark.table("ods.ods_order_detail") \ .filter(col("dt") == "20240101") verification_result = VerificationSuite() \ .onData(data) \ .addCheck( Check(CheckLevel.Error, "订单明细表深度质量检查") .isComplete("order_id") .isUnique("order_id") .isNonNegative("order_amount") .isContainedIn("order_status", ["待支付", "已支付", "已取消", "已退款", "已发货"]) .hasSize(lambda size: size >= 1000000) ) \ .run() verification_result_df = VerificationResult.successMetricsAsDataFrame(spark, verification_result) verification_result_df.show(truncate=False)这段代码里,isComplete对应完整性,isUnique对应唯一性,isNonNegative对应金额合法性,isContainedIn对应枚举值检查,hasSize对应行数下限。Deequ有个很大的优势是可以输出成功指标(success metrics),比如唯一率、缺失率、通过与否,这些指标非常适合存入质量结果表做趋势分析。
4.4 用Apache Griffin做跨源对账
还有一个很典型的场景:业务库MySQL里的订单表和数仓Hive里的订单表,每天数据量是否一致?这个用Soda Core可以分别查两个表再对比,但过程比较繁琐,专业一点的方式是用GoGriffin。Griffin的准确度(Accuracy)检查专门用来做这种跨数据源的对账,它会计算两张表按照主键关联后不匹配的记录数。
Griffin的规则配置虽然用JSON写,但逻辑更像写SQL。核心配置是这样的形式:
{ "name": "order_accuracy_check", "data.source": "batch", "source": { "data.connector": { "type": "hive", "config": { "database": "ods", "table.name": "order_detail" } } }, "target": { "data.connector": { "type": "jdbc", "config": { "database": "order_biz", "table.name": "order_detail", "user": "quality_ro", "password": "your_password" } } }, "evaluate.rule": { "rules": [ { "rule.dsl.type": "griffin-dsl", "rule": "src.order_id = target.order_id and src.order_amount = target.order_amount and src.order_status = target.order_status" } ] } }这条规则的意思是:按照order_id关联后,比较两边的order_amount和order_status是否一致,不一致的就是质量问题。Griffin会跑一个Spark任务做全量比对,输出一个accuracy的指标。说实话,Griffin的安装和配置成本在四个工具里是最高的,如果你们团队暂时没有跨源对账的硬需求,可以放在第二阶段再引入。
5. 调度编排、结果入库与告警通知
5.1 调度任务设计
调度是整个平台的中枢。我用DolphinScheduler维护了一套工作流,每天凌晨5点开始执行质量校验。工作流设计成四个步骤:
- 等待数仓ETL任务跑完(通过依赖检查判断,比如先检查
ods_order_detail分区是否存在且表内行数大于0); - 执行Soda Core检查;
- 执行Deequ深度校验;
- 汇总结果,推送日质量报告。
DolphinScheduler里创建shell任务,命令示例如下:
#!/bin/bash source /opt/soda/venv/bin/activate cd /opt/quality_monitor soda scan -d order_dw -c configuration.yml checks.yml --output-dir results/$(date +%Y%m%d)调度周期设置成每天一次,失败重试一次。如果检查结果异常,DolphinScheduler本身也有告警插件,但我们不走那个,因为Soda Core的退出码更细,按检查结果分级告警更准确。Soda Core执行完检查后,如果有一个error级别的检查失败,进程的退出码就是非零,我们可以在shell脚本里捕获这个状态再触发告警。
5.2 结果入库脚本
Soda Core生成的JSON结果,需要解析后写入PostgreSQL的质量结果表。我写了一个简单的Python脚本做解析和入库,核心逻辑如下:
import json import glob import psycopg2 from datetime import datetime def parse_soda_result(file_path): with open(file_path, "r", encoding="utf-8") as f: data = json.load(f) rows = [] for table_result in data.get("checks", []): rows.append(( data.get("metadata", {}).get("execution_date", "")[:10], table_result.get("table", "").split(".")[-1], table_result.get("name", ""), table_result.get("outcome", "Unknown"), table_result.get("metric", ""), table_result.get("value", None), str(table_result.get("threshold", "")), datetime.now() )) return rows def insert_rows(rows): conn = psycopg2.connect( host="192.168.1.10", port=5432, dbname="quality_meta", user="quality_writer", password="your_password" ) cur = conn.cursor() insert_sql = """ INSERT INTO data_quality_check_result (exec_date, table_name, check_name, check_status, metric_name, metric_value, threshold, create_time) VALUES (%s, %s, %s, %s, %s, %s, %s, %s) """ cur.executemany(insert_sql, rows) conn.commit() cur.close() conn.close()注意一个细节:写库账号我用quality_writer,跟ODS数据源账号分开。生产环境读写权限分离是基本要求。
5.3 告警推送:钉钉、企业微信、邮件
告警是数据质量平台的“最后一公里”,处理不好前面全白搭。我遇到过最尴尬的情况是:规则确实检查出来了问题,但告警发到了群聊里被其他消息刷掉,最后问题还是靠业务方发现的。
我现在用的是企业微信机器人Webhook。Python脚本实现简单,把结果拼成Markdown消息推送到指定群里:
import requests import json import sys def send_wechat_alert(content): webhook_url = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=YOUR_KEY" payload = { "msgtype": "markdown", "markdown": { "content": content } } resp = requests.post(webhook_url, json=payload) if resp.status_code != 200: print("webhook send failed:", resp.text) if __name__ == "__main__": title = sys.argv[1] detail = sys.argv[2] msg = f"## 数据质量告警\n**告警标题:** {title}\n**详情:** {detail}\n<@all>" send_wechat_alert(msg)告警策略上,我建议做分级处理:
- error级别:当检查项失败时,立刻推送群消息,并@值班人。
- warning级别:只写入结果库和日报,不推送群消息,避免打扰。
- 同规则连续失败超过3天:升级为电话告警或高优先级工单。
5.4 调度脚本整合
把上述串起来,一个完整的调度脚本大概是这样的:
#!/bin/bash export PATH=/usr/local/bin:/opt/soda/venv/bin:$PATH TODAY=$(date +%Y%m%d) cd /opt/quality_monitor # Step1: run soda soda scan -d order_dw -c configuration.yml checks.yml --output-dir results/$TODAY SODA_EXIT=$? # Step2: parse and store python3 result_parser.py results/$TODAY # Step3: alert if needed if [ $SODA_EXIT -ne 0 ]; then python3 alert.py "订单表数据质量校验未通过" "详情请看质量看板,执行日期: $TODAY" fi这个脚本可以直接放进crontab,也可以挂在DolphinScheduler下。我建议在DolphinScheduler里配置,因为它自带日志、失败重跑和依赖管理,排查问题方便得多。
6. 可视化看板与日质量报告
6.1 Grafana看板设计
数据质量结果全部落库之后,我第一时间去配了Grafana看板。为什么要看板?因为“检查通过”和“数据质量好”是两回事。质量指标必须是可趋势化的,需要看到连续几天的通过率变化、波动幅度,才能提前发现苗头。
我在Grafana里创建了几个核心面板:
- 今日整体通过率:统计当天所有检查项中通过数量占总数的比例。
- 按表维度检查结果:展示每张表的error数量、warning数量,用表格形式罗列,方便晨会快速过。
- 关键指标趋势:比如订单表的行数波动率、金额非空率连续30天趋势,用折线图展示。
- 失败规则明细:列出最近一天失败的具体规则和失败值,直接定位问题。
Grafana接入PostgreSQL数据源后,查询语句可以这样写:
SELECT exec_date, table_name, count(*) FILTER (WHERE check_status = 'Pass') AS pass_count, count(*) FILTER (WHERE check_status = 'Fail') AS fail_count FROM data_quality_check_result WHERE exec_date >= CURRENT_DATE - INTERVAL '30 days' GROUP BY exec_date, table_name ORDER BY exec_date DESC;6.2 日质量报告与周会抓手
看板适合实时看,日报适合沉淀和复盘。我每天早上9点用Python脚本生成一份前一日的质量日报,内容包括:核心表通过率、环比波动、连续失败规则列表、需要人工确认的异常项。发送到数据团队群里。
这个日报不要写太复杂,核心就是让负责人30秒能扫完。生成日报可以用简单的SQL统计,再把结果用模板拼装成Markdown,通过企业微信机器人发出去。格式大致如下:
SELECT table_name, check_name, check_status, metric_value, threshold FROM data_quality_check_result WHERE exec_date = CURRENT_DATE - 1 AND check_status = 'Fail' ORDER BY table_name, check_name;如果查询结果为空,就只在日报里写“昨日所有检查通过”,很简单。
6.3 质量分的计算
日报和看板都不能只展示通过/失败,最好有一个综合“质量分”来评价整体健康度。我的做法是给每种检查项设置权重,通过一个SQL直接算出当天的质量得分:
SELECT ROUND( 100.0 * SUM(CASE WHEN check_status = 'Pass' THEN weight ELSE 0 END) / SUM(weight), 2 ) AS quality_score FROM data_quality_check_result JOIN quality_rule_weight ON data_quality_check_result.check_name = quality_rule_weight.check_name WHERE exec_date = CURRENT_DATE;质量分不是为了KPI,而是为了让大家能直观看到整体趋势。如果某天质量分从99掉到80,说明肯定有环节出了问题,值得让人去翻一下。
7. 常见问题排查与避坑实录
7.1 Soda Core连接数据源失败
最常遇到的问题是执行soda scan时报无法连接数据源,但用命令行工具却能连上。这类问题90%是以下原因之一:
- 驱动未安装:Soda Core必须安装对应数据源的插件包,比如
pip install soda-core-mysql、soda-core-postgres,只装soda-core主包是不够的。 - 网络隔离:数据质量平台部署机器到数据源之间网络不通,检查防火墙和VPC访问白名单。
- 只读账号权限不足:某些数据源驱动会先读取库表元数据,如果账号连information_schema的读取权限都没有,也会报连接失败。
排查命令很简单,先用telnet测试网络连通性,再确认驱动包。
7.2 Deequ在Spark 3.x上运行报错
Deequ对Spark版本比较敏感,早期版本在Spark 3.2以上容易碰到方法签名不兼容的问题。我遇到过的报错是NoSuchMethodError: org.apache.spark.sql.catalyst.plans.logical.LogicalPlan,后来排查发现是Maven里引入了多个版本的spark-sql依赖冲突。
解决方式是确保只存在一个Spark版本依赖,并且在提交Spark任务时使用集群的spark-submit,不要另起Python进程。如果用的是Python版Deequ,建议把Spark版本和Deequ版本锁定,用下面的命令在本地验证兼容性:
pip show deequ spark-submit --version7.3 Griffin任务一直pending不执行
Griffin任务提交到YARN后一直处于等待状态,这个坑我印象很深。主要原因通常是资源队列紧张,或者Griffin服务所在的节点跟YARN集群不在同一个Kerberos环境里。排查时先看YARN的资源剩余量:
yarn application -list yarn node -list -all如果资源没问题,再看Griffin配置里的Spark执行参数。我建议把Executor内存和核数调小一点,让任务能快速进入调度,不要一上来就申请大资源。
7.4 告警轰炸与误报
告警轰炸是数据质量监控用起来最让人头疼的问题。规则一旦配置过严,每天凌晨定时任务没跑完或者稍微延迟几秒,就会刷出一堆无用告警,导致没人再看告警消息。我自己的策略是:
- 新规则先以
warning级别灰度跑一周,观察通过率再决定是否升到error; - 相同的规则、相同的表连续失败,只推送第一条告警,并在日报中汇总后续失败情况;
- 告警消息里必须带上
exec_date、规则名称、失败数值和阈值,让收到消息的人不用再开电脑查。
7.5 数据新鲜度检查的时区问题
检查数据新鲜度时,我踩过一个特别隐蔽的坑。数仓里分区时间用的是业务时区,服务器用的是UTC时间,直接比较max(order_time)和now()会相差8小时,导致每天早上都误报“数据延迟”。后来我在写好检查规则时把时间比较统一转为业务时区,或者干脆在SQL中把两个时间都转成时间戳再比较。
7.6 结果表数据无限膨胀
质量校验结果每天都会产生大量数据,而且检查项越多数据量越大。如果不做清理,一年下来结果表会很大,影响查询性能。我的做法是:
- 结果表按
exec_date做分区,如果用的是PostgreSQL,可以把老数据定期迁移到归档表或者直接清理。 - 只保留最近90天的完整结果,超过90天的按天聚合后存入汇总表,只保留通过率、失败数量等指标。
这一步看似不起眼,但半年后你回头看,会发现当初做了数据清理决策太正确了。
8. 从0到1落地的一点个人建议
这套平台我前后迭代了三版。第一版全部用crontab定时跑脚本,规则写死在Python里,结果打到日志文件,要查的时候用grep翻。第二版引入了Soda Core,规则用YAML管理,结果入库,有了第一个看板和告警。第三版才真正接入了DolphinScheduler和Deequ,形成了完整的工具链闭环。
如果让我从头再搭一次,我会给团队一个更明确的阶段性目标:第一周只做两张核心表的完整性和唯一性巡检,用crontab跑通告警链路;第二周把结果入库、看板做出来;第三周再上Deequ深度校验和跨源对账;调度统一迁移放到一个月后。不要想着一步到位把所有表都监控起来,先让两张核心表稳定运行一个月,用实际效果说服团队和领导,再慢慢推广到更多表和更多规则。
数据质量监控平台的价值不是上线当天就能看出来的,需要连续运行几周,积累足够多的趋势数据后,才能体现出“提前发现异常”的意义。我个人最后的体会是:不要过度迷信某一个工具,也不要一开始就追求大而全的架构,先把最小闭环跑起来,再把链路补完整,这才是最稳妥的路线。