大数据数据质量监控平台搭建:基于开源工具的一站式解决方案
做大数据平台三年多,我踩过最深的坑不是集群性能,不是数据倾斜,而是"数据错了没人知道"。数仓里跑了几十张表,上游一个字段类型变更、接口某个枚举值调整,下游报表算出来的数悄无声息就偏了。业务方拿着数去决策,等发现的时候,损失已经造成。
后来我下定决心,把数据质量监控当成一个独立平台来建设,而不是零散地写脚本、定时跑 SQL。今天这篇就完整复盘一下,我如何基于 Apache Griffin、Apache Atlas、Apache Airflow 和一套灵活的规则引擎,搭出一个覆盖完整性、准确性、一致性、唯一性、及时性五个维度的开源数据质量监控平台。文章会从整体架构设计、核心组件选型、实操实现到避坑指南都过一遍,适合正在做数仓治理、数据中台建设,或者对数据质量一脸迷茫的开发人员参考。
1. 理解数据质量:先搞清楚要监控什么
在动手选工具之前,我建议所有人都先把"数据质量"这件事本身想清楚。否则工具选得再花哨,也不知道该监控哪些东西。
1.1 数据质量坑在哪:五个维度说透
所谓数据质量监控,本质是回答一个问题:数据在流转过程中,有没有出现不符合预期的"异常"。我把它拆成五个可量化的维度,这也是目前业界比较公认的划分方式。
完整性:该有的数据有没有缺失。比如用户表应该有 1000 万条记录,实际只有 980 万;订单表某一天的分区为空;某个关键字段如 user_id 有大量 NULL。
准确性:数据值本身对不对。比如订单金额出现了负数、年龄字段出现了 200 岁、身份证号位数不对。
一致性:同一份数据在多个系统中表达是否统一。最典型的就是数仓中订单金额的单位,源系统用"分",数仓用"元",换算错了之后两边的汇总对不上。
唯一性:主键是否唯一。比如用户维度表 user_id 出现重复,导致 join 的时候数据翻倍。
及时性:数据能否在预期时间内到达。比如 T+1 的数据任务应该在凌晨 6 点前产出,但是某天延迟到 9 点才跑完,下游报表就会变成"昨天的数据"。
这五个维度并不是互相独立的。比如源系统某个字段的枚举值更新了,可能同时引发一致性和准确性问题。所以监控规则的设计不能只盯着一个维度,需要组合起来看。
1.2 为什么不能靠人工查数:质量监控的被动与主动
很多团队一开始都是靠"人工发现"——业务方反馈数据不对了,数据工程师再去查。但数据量一大,这种模式基本失效。我见过一个典型场景:数仓里挂了 2000 多张表,每天跑几万个任务,数据工程师能盯得过来的,只有业务方天天看的那几十张核心表。剩下那些偶发问题,往往要等月底对账的时候才会暴露。
更麻烦的是"脏数据会传染"。A 表质量没问题,但 B 表 join 的时候用了错误的关联键,导致 B 表数据翻倍,而 C 表又依赖 B 表……等发现问题的时候,一整条链路的数据都要重刷。这比上游源系统出问题还要耗时,因为数据血缘复杂,排查链路极长。
所以数据质量监控平台必须"主动"运行:周期性自动跑规则,一旦异常就立刻告警,把问题拦截在报表产出之前,而不是等下游业务方来投诉。这也是这套平台建设的核心目标。
2. 开源工具选型:为什么我选 Griffin + Atlas + Airflow
先说明一下,我不推荐从零开始写一套数据质量校验系统。原因很简单:数据质量监控涉及规则定义、任务调度、指标存储、血缘分析、告警通知,每个环节都有成熟方案,自己写看起来灵活,实际上维护成本极高。
2.1 面向场景的选型思路
我的选型约束条件很清晰:
- 必须开源且社区活跃,不能是单点维护的个人项目;
- 能跑在已有的 Hadoop/Spark 集群上,不额外引入重组件;
- 规则要灵活,既能做表级校验,也能做字段级、跨表级的校验;
- 要能与数仓的调度体系打通,最好能自动获取血缘关系。
这些条件筛选下来,核心组件就锁定在三个:Apache Griffin 负责指标计算、Apache Atlas 负责血缘和元数据、Apache Airflow 负责周期性调度。辅助组件包括一个规则配置库(我用的 MySQL)和一个可视化看板(当时用的 Grafana,后来也试过 Superset)。
2.2 三个核心组件的作用边界
Apache Griffin:这是 LinkedIn 开源的数据质量工具,核心能力是定义"数据质量测度",然后将测度翻译成 Spark 任务执行。它支持准确性、完整性、一致性等常见校验类型,结果以 JSON 的格式写入 HDFS,同时可以对接 ES 做查询。Griffin 是这套平台的计算引擎,也是最容易踩坑的部分,后面我详细讲。
Apache Atlas:Atlas 是 Hadoop 生态的元数据管理与数据血缘工具。我选它主要不是为了"好看",而是为了解决问题定位:某张表数据异常,能不能立刻顺着血缘找到上游的影响范围。Atlas 支持通过 Hook 自动采集 Hive、Spark 的元数据和血缘信息,配合平台做"血缘追溯"。
Apache Airflow:调度器。数据质量监控任务需要按业务节奏周期性执行,Airflow 用 DAG 定义依赖关系,特别适合数仓的 T+1 批处理场景。Griffin 的任务可以封装成 Airflow 的一个 operator,到点自动触发,失败自动重试加告警。
一句话总结:Griffin 负责"算"质量,Atlas 负责"管"元数据和血缘,Airflow 负责"跑"定时任务。三者结合,再加上告警通知,就形成了一个完整的闭环。
| 组件 | 定位 | 选型理由 | 替代方案 |
|---|---|---|---|
| Apache Griffin | 数据质量指标计算 | 支持多种测度,原生对接 Spark,结果可写入 ES/HDFS | Deequ(基于 Spark 的库,需要自己包调度) |
| Apache Atlas | 元数据与血缘管理 | 原生支持 Hive/Spark Hook,血缘自动采集 | DataHub(偏数据目录,血缘能力可参考) |
| Apache Airflow | 任务调度 | 生态成熟,DAG 灵活,可自定义 operator | DolphinScheduler(国产,界面友好,但算子自定义略繁琐) |
| Grafana/Superset | 展示和告警 | 接 ES/MySQL,快速出图、配置告警 | 也可以直接看 Griffin 自带的 UI |
3. 平台总体架构设计
工具选完,接下来是架构设计。我把整个平台分成五层,每一层职责单一,这样后续扩展和排障都比较清晰。
3.1 五层架构详解
从上到下分别是:
采集层:对接元数据源(Hive、Kafka、业务库 binlog),负责把要监控的表、字段、分区信息采集到平台自己的元数据仓库里。这一层我直接用 Atlas 的 Hook 自动采集,不需要写太多代码。
规则层:这是平台的核心。用户可以通过配置中心定义数据质量规则,包括规则类型(完整性、准确性、唯一性等)、适用表/字段、阈值、生效时间、告警级别。规则配置存在 MySQL 中,Griffin 执行任务时会从配置中心读取规则并翻译成 Spark SQL。
调度层:Airflow 按照配置的周期触发质量校验任务。每种规则生成一个可执行的 Spark 任务,任务执行结果回写到结果表。
存储与计算层:底层依赖已有的 Hadoop 集群,Spark 负责执行校验计算,校验结果落到 HDFS/ES,方便快速查询。
展示层:Grafana 接 ES/MySQL 数据源,展示表/字段级别的质量得分、趋势图、告警事件;Atlas UI 用来查看血缘关系,定位问题影响范围。
3.2 为什么必须引入数据血缘
很多人觉得血缘是"锦上添花",实际用起来才发现它是"雪中送炭"。数据质量告警之后,下一步就是"影响分析":这张数仓表的数据不对,下游有哪些指标、哪些报表会受影响?人工梳理根本来不及,尤其数仓层级深的时候,A→B→C→D 之间跨了五层。
Atlas 通过解析 Hive/Spark 的执行日志,自动建立表与表、字段与字段之间的血缘关系。比如我发现"dws_order_daily"表数据总量异常,在 Atlas 里点开这张表,就能看到它依赖哪些 dwd 层表,以及它又供给了哪些应用层表。配合数据质量平台,告警之后可以立刻评估影响范围,决策是"继续跑"还是"紧急修数"。
3.3 规则与任务的映射关系
数据质量规则不是一堆散落的配置,而是有明确结构的。我设计的时候,把一条规则拆成四个属性:
- scope:校验范围。是表级,还是字段级,还是跨表 join 校验。
- metric:校验指标。比如 null_count、distinct_count、total_count、duplicate_count、max_value、期望平均值。
- condition:过滤条件。比如只校验 where dt = '2025-01-01' 的分区。
- threshold:阈值阈值。比如 null 率不能超过 1%,或总量波动不能超过 ±5%。
Griffin 原生支持 accuracy(准确性)、completeness(完整性)、distinctness(唯一性)等测度,本质上就是一组预定义的 Spark SQL template。实际使用中,我用得最多的是"完整性 + 唯一性 + 自定义准确性 SQL"组合,后面第 4 节会给出具体配置示例。
4. 平台部署与规则实现:从零到上线
这一节是实操重点,我会按照我实际搭建的顺序来写。需要说明的是,具体版本号并不绝对,关键是理解每步在做什么。
4.1 环境准备:版本匹配是最大的坑
Griffin 对 Spark 版本比较敏感,这是整个搭建过程中最容易让人崩溃的地方。我的环境是:
- Hadoop 3.1.1 / Hive 3.1.2
- Spark 2.4.8(Griffin 官方支持的版本之一)
- Apache Griffin 0.6.0
- Apache Atlas 2.2.0
- Airflow 2.x(我用的 2.5.1)
- MySQL 8.0(规则配置库)
- Elasticsearch 7.x(指标存储与查询,用于 Grafana 展示)
注意:Griffin 0.6.0 官方适配的是 Spark 2.4.x,如果你用的是 Spark 3.x,编译 Griffin 源码时要额外处理依赖冲突。我自己试过 Spark 3.1.1,踩了不少坑,后面第 5 节会细说。
这个环境里还有一个关键点:Hive 和 Spark 的元数据要打通。Griffin 生成的 Spark 任务要能直接读 Hive 表,所以 Spark 的 hive-site.xml 必须指向 Hive 的 Metastore,否则会报"Table not found"。
4.2 数据源接入与元数据采集
数据源接入这块,我用了 Atlas 的 Hive Hook 来做元数据自动采集。部署方式:
- 把 atlas-application.properties 配置好,指向 Atlas 服务端;
- 将 Atlas Hook 的 jar 包放到 Hive 的 auxlib 目录;
- 重启 HiveServer2 之后,所有 Hive 的 DDL、Query 都会自动上报到 Atlas,表结构和血缘自动就有了。
注意一点:如果表数量非常多,全量采集一次 Atlas 会比较慢。建议先在 Atlas 里只采集核心库和核心表,跑通之后再放开。
4.3 编写第一条数据质量规则
规则定义我用的是 JSON 配置,存在 MySQL 里。下面是"订单表完整性校验"的一个简化示例:
{ "job_name": "dq_order_daily_completeness", "data_source": "hive", "table_name": "dwd_order_daily", "measure_type": "completeness", "rule": { "target_field": "order_id", "null_threshold": 0.01 }, "schedule": { "cron": "0 30 2 * * ?", "timezone": "Asia/Shanghai" }, "alert": { "level": "high", "channels": ["webhook", "email"] } }这条规则的含义是:每天凌晨 2:30 执行,检查 dwd_order_daily 表中 order_id 字段的 NULL 率,如果超过 1% 就触发告警。
在 Griffin 里,这个配置最后会被翻译成 Spark 任务执行,核心逻辑类似:
SELECT COUNT(*) AS total_count, SUM(CASE WHEN order_id IS NULL THEN 1 ELSE 0 END) AS null_count FROM dwd_order_daily WHERE dt = '2025-01-01'4.4 Airflow 调度集成与告警闭环
Griffin 本身没有调度能力,需要靠外部触发。我封装了一个 Airflow operator,把上述 JSON 配置变成可执行的 Spark 提交命令:
from airflow.models import BaseOperator from airflow.utils.decorators import apply_defaults import subprocess class GriffinDQOperator(BaseOperator): @apply_defaults def __init__(self, griffin_job_name, spark_submit_cmd, *args, **kwargs): super().__init__(*args, **kwargs) self.griffin_job_name = griffin_job_name self.spark_submit_cmd = spark_submit_cmd def execute(self, context): self.log.info(f"Submitting Griffin job: {self.griffin_job_name}") process = subprocess.Popen( self.spark_submit_cmd, shell=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT ) exit_code = process.wait() if exit_code != 0: raise Exception(f"Griffin job failed with exit code {exit_code}") self.log.info("Griffin job completed successfully")Airflow DAG 里这样定义调度:
from airflow import DAG from airflow.utils.dates import days_ago default_args = { 'owner': 'data_quality', 'retries': 2, 'retry_delay': 300 } with DAG( dag_id='dq_daily_check', default_args=default_args, start_date=days_ago(1), schedule_interval='30 2 * * *', catchup=False ) as dag: dq_order_daily = GriffinDQOperator( task_id='dq_order_daily', griffin_job_name='dq_order_daily_completeness', spark_submit_cmd='./spark-submit --class org.apache.griffin...' )告警通知这一环,我最初是直接让 Griffin 写结果到 ES,Grafana 配置阈值告警,通过 Webhook 推到企业微信群。后来发现,更合理的做法是让 Airflow 失败时直接触发告警——因为 Airflow 重试和失败机制更成熟,而且可以带上执行上下文信息。所以最终方案是:规则校验结果指标走 Grafana 展示,任务失败告警走 Airflow。
5. 平台落地过程中的典型问题与排查实录
说实话,这套平台真正的价值是在一次次踩坑之后才体现出来的。下面这些问题是几乎每个搭建者都会遇到的,我按频率排序。
5.1 Spark 版本与依赖冲突
Griffin 官方包编译时用的是 Spark 2.4,如果你的集群是 Spark 3.x,直接拿来跑大概率会报类似的依赖错误:
java.lang.NoSuchMethodError: org.apache.spark.sql.execution.datasources.FileFormatWriter$.write这是 Spark 内部 API 变化引起的,不是配置问题。我当时的解决办法是:源码编译 Griffin 0.6.0,在编译命令里指定 Spark 版本:
mvn clean package -DskipTests -Pspark-2.4如果你的集群确实是 Spark 3,需要修改 Griffin 源码中的依赖版本,重新交叉编译,这个工作量和风险都比较高。我个人建议,如果不是非用 Spark 3 不可,就让 Griffin 跑 Spark 2.4,单独部署一套"质量校验计算集群",跟主集群做资源隔离,互不干扰。
5.2 指标膨胀与 HDFS 小文件问题
Griffin 每次跑任务,都会往 HDFS 写一份结果 JSON。如果规则很多、频率又高,小文件问题会让 NameNode 压力越来越大。最简单有效的方式是定期把结果文件压缩合并:
hdfs dfs -ls /griffin/dq_result/ # 按天归档 hdfs dfs -mv /griffin/dq_result/*.json /griffin/archive/2025-01-01/或者更彻底一点,直接把 Kylin 的存储思路借过来——用 Spark 定期读结果目录,重新写成分区表的 ORC 文件。我在第 3 个月的时候做了一次改造,把结果统一落到 Hive 的外部表,按 dt 分区,再通过 Presto 供 Grafana 查询,性能和稳定性都好了很多。
5.3 规则误报与阈值制定
阈值设置不合理,是数据质量平台最容易被业务方"骂"的地方。阈值太松,异常发现不了;太紧,天天误报,大家就麻木了。
我的经验是:阈值不能拍脑袋,要用历史数据来确定。先跑两周的"只记录不告警"模式,把每个指标的基线数据存下来,比如总行数、NULL 率、重复率,然后基于均值 ± 3 倍标准差做动态阈值,或者取 P95/P99 分位数做静态阈值。比如订单日表的行数,过去 30 天波动率在 ±3% 以内,那阈值就设 5%,留出一定缓冲。
5.4 Atlas 血缘信息不全的问题
Atlas 的表面血缘靠 Hive 的 DDL 和 Query Hook 采集,但实际运行中,很多人会发现血缘只有表级,没有字段级,或者某些临时表没有血缘。排查后发现原因有两个:
- 临时表(tmp_xxx)太多,Hook 采集的信息大量冗余;
- 某些 SQL 是通过 Spark SQL 跑的,没有配置 Spark Hook。
建议在 Atlas 里做表名白名单过滤,只采集正式库的表;同时给 Spark 也配上 Atlas Hook,保证 Spark 作业产生的血缘也能采集到。
血缘采集的及时性也要注意:Atlas Hook 是异步上报的,任务跑完之后血缘可能需要几分钟才能体现在 UI 里,这属于正常现象,不用焦虑。
5.5 权限管理与多租户
平台用起来之后,你会发现不同团队都想配置自己的监控规则。如果所有规则都放在一个配置中心,会乱成一团。我的做法是引入简单的"命名空间"概念:
- 每条规则属于一个项目(project),项目有 owner;
- 指标结果表按 project 分区,查数据和看板按项目隔离;
- owner 只能看到自己项目下的规则和告警记录。
权限这块我直接用 MySQL 的表级权限 + Airflow 的 DAG owner 权限来控制,没有引入额外的权限组件。如果你们的团队规模超过 20 人,再考虑引入统一权限体系。
6. 效果复盘与后续扩展
平台上线三个月之后,我的直接感受是:告警数量从最初的每天几十条,降到了每天三五条,而且剩下的基本都是真实问题,不是误报。这种"信任感"是最大的回报,业务方开始愿意主动看质量看板,而不是出了问题再来找数仓。
从指标上看,平台覆盖了 300 多张核心表的 1500 多条规则,每天跑 200 多个质量校验任务,平均每个任务耗时 3 分钟以内,对集群资源的占用基本可以忽略。问题发现平均提前量大概在 4 到 6 小时,比人工发现快了不止一个量级。
6.1 还能怎么延伸:从监控走向治理
监控只是第一步,纯监控并不能修复问题。后续扩展我建议往两个方向走:
自动修复:对于少数能够预先定义修复逻辑的问题(比如异常数据回退、重新拉取源数据),可以在告警触发后自动跑一个修复 DAG,而不是等人来处理。
质量分与奖惩机制:把表级别的质量得分纳入数仓开发流程。质量分低于某个阈值的表,不允许发布到生产;开发了质量监控规则的表,可以有相应的资源倾斜。这样从制度上让"质量"和"开发"绑定在一起。
我个人觉得,数据质量平台的价值,不在于你能不能写出复杂的监控规则,而在于它能不能真正融进团队的日常开发节奏里。如果每个表发布之前都必须配套质量监控规则,每个告警都必须有人认领和处理,那么这个平台就不只是一个工具,而是一套质量文化的基础设施。说到底,好数据不是查出来的,是大家都把它当回事之后,一点一点养出来的。