news 2026/9/9 19:29:43

用开源工具搭建一站式数据质量监控平台实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
用开源工具搭建一站式数据质量监控平台实战

做数仓的兄弟应该都经历过这种深夜打来的电话:报表跑完了,但一看数据明显不对,订单量少了三分之一,查了半天定位到上游同步任务凌晨挂了,重跑之后才发现脏数据已经污染了下游。数据量越大、链路越长,这种问题就越难靠人肉巡检兜住。所以这两年我一直在推进一件事:用开源工具拼一套数据质量监控平台,让规则自动跑、结果自动存、异常自动喊人。

这篇文章就把这套方案的完整搭建过程写出来。文章不侧重某个商业产品,而是围绕“一站式”这个目标,把开源生态里几个能打的数据质量工具按场景组合起来:离线大表用Deequ做深度校验,业务表快速接入用Soda Core做轻量检查,跨源数据对账用Apache Griffin,调度统一走DolphinScheduler,结果统一入库,告警统一推送,展示统一做Grafana看板。适合数仓工程师、数据平台开发、数据治理方向的同学参考,也适合团队刚准备做数据质量、还没想清楚怎么落地的朋友拿去直接用。

1. 先从业务痛点说起:为什么要自建数据质量平台

1.1 数据质量问题的真实代价

很多人一开始觉得“数据质量”是个很虚的概念,直到线上出了事故才意识到它的分量。最常见的场景是ETL链路里某个字段映射错了,或者源系统数据库扩容导致同步任务凌晨挂了然后自动重跑,重跑出来的结果跟正常业务时间对不上。这类问题往往不是立刻引爆的,而是等下游报表跑完、业务方早上拿到数据时才发现异常,这时候已经过去了两三个小时。

数据质量问题带来的直接代价不只是“数据不准”,还有信任崩塌。业务方会因为一次错误数据对整个数仓产生怀疑,之后每一张报表都要反复确认数据来源。这种信任一旦丢失,数据平台的生存空间会被大大压缩。更麻烦的是,当数据量从千万级涨到亿级、表数量涨到几千张时,纯靠人工去抽数、查数、盯数根本不可能,必须有系统化的工具在无人值守状态下持续运行。

1.2 平台需要解决的核心问题

我整理了一下,一个能真正落地的数据质量监控平台至少要解决四件事:

  • 规则可配置:不同表、不同字段的质量标准不一样,订单表要求订单号不能为空,用户表要求邮箱格式正确,日活表要求行数波动不超过10%。规则必须能按表、按字段灵活配置,最好能用声明式写法,不用开发写一堆Java代码。
  • 调度可管理:每天凌晨离线任务跑完后自动触发质量校验,校验任务要能被统一调度、重跑、补数,还能看到每个任务跑了多久、是否成功。
  • 结果可追溯:某一天的数据质量不合格,要能查清楚是哪张表、哪个字段、哪个规则出了问题,当时的校验结果是什么样的,不能等业务方投诉了才去翻日志。
  • 告警可送达:规则失败后要第一时间通知到责任人,而且不能一告警就刷屏,要有分级和分组策略。

如果用商业产品,比如Informatica、Collibra,功能确实全,但价格不是一般团队能承受的。用开源工具自己拼,成本低、可控性强,缺点是得自己踩坑。我下面要讲的方案,就是围绕这四件事去搭建的。

1.3 为什么选择开源工具而不是自研

不少团队一开始想自研数据质量平台,理由是“我们数据量特殊,只有自己写的才符合需求”。这个思路我不反对,但强烈建议先站在开源工具的肩膀上起步。原因很简单:数据质量的核心在于规则引擎和指标计算,这块的通用性远比你想象的高。完整性、唯一性、及时性、准确性,不管什么行业都逃不开这几个维度,开源工具已经沉淀了很多经验。

自研需要投入研发人力做规则解析、做任务调度、做结果存储、做告警通知,这些都属于“轮子”,而且做出来的稳定性大概率不如经过社区考验的成熟工具。我建议的路线是:先用开源工具把核心链路跑通,如果后续有非常个性化的需求,再在开源工具之上做扩展开发,而不是从头写一套。

2. 平台整体架构与开源工具选型

2.1 常见开源数据质量工具对比

目前社区里主流的开源数据质量工具主要有四个,我把它们放在一起对比过,各自的侧重点很不一样。

工具底层引擎擅长场景规则定义方式社区活跃度
Apache GriffinSpark跨源数据对账、批流一体质量监控JSON配置中低
DeequSpark离线大表多维度校验、数据画像Scala/Python API中高
Soda CorePython业务库/数仓表快速接入、轻量级检查YAML声明式
Great ExpectationsPython数据管道测试、数据文档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_iduser_idorder_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有两个级别:errorwarning。我个人建议核心表全部用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_amountorder_status是否一致,不一致的就是质量问题。Griffin会跑一个Spark任务做全量比对,输出一个accuracy的指标。说实话,Griffin的安装和配置成本在四个工具里是最高的,如果你们团队暂时没有跨源对账的硬需求,可以放在第二阶段再引入。

5. 调度编排、结果入库与告警通知

5.1 调度任务设计

调度是整个平台的中枢。我用DolphinScheduler维护了一套工作流,每天凌晨5点开始执行质量校验。工作流设计成四个步骤:

  1. 等待数仓ETL任务跑完(通过依赖检查判断,比如先检查ods_order_detail分区是否存在且表内行数大于0);
  2. 执行Soda Core检查;
  3. 执行Deequ深度校验;
  4. 汇总结果,推送日质量报告。

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-mysqlsoda-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 --version

7.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深度校验和跨源对账;调度统一迁移放到一个月后。不要想着一步到位把所有表都监控起来,先让两张核心表稳定运行一个月,用实际效果说服团队和领导,再慢慢推广到更多表和更多规则。

数据质量监控平台的价值不是上线当天就能看出来的,需要连续运行几周,积累足够多的趋势数据后,才能体现出“提前发现异常”的意义。我个人最后的体会是:不要过度迷信某一个工具,也不要一开始就追求大而全的架构,先把最小闭环跑起来,再把链路补完整,这才是最稳妥的路线。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/9 19:27:25

Kafka消息可靠性全链路保障:从生产端到消费端的配置与排查实战

1. Kafka消息可靠性的核心问题与设计思路做大数据的人没几个没被Kafka折磨过。我最早接触Kafka的时候&#xff0c;以为它默认配置就能放心用&#xff0c;结果线上数据一丢就是几万条&#xff0c;排查到半夜才发现&#xff1a;生产端acks没配、Broker端副本数不够、消费端位移自…

作者头像 李华
网站建设 2026/9/9 19:25:26

5 分钟跑通 Agentic 测试:从最小用例到快照防回归

5 分钟跑通 Agentic 测试&#xff1a;从最小用例到快照防回归 【免费下载链接】agentic Your API ⇒ Paid MCP. Instantly. 项目地址: https://gitcode.com/GitHub_Trending/ag/agentic Agentic 测试不需要测试理论基础&#xff1a;读完全文&#xff0c;你能独立在这个仓…

作者头像 李华
网站建设 2026/9/9 19:23:01

嵌入式Linux下librtmp交叉编译实战:从依赖到部署

简介&#xff1a;librtmp库作为rtmpdump工具的核心组件&#xff0c;长期用于RTMP协议处理与实时流媒体开发。该librtmp库实测可在Linux环境下完成交叉编译&#xff0c;面向需要在x86开发机上为ARM等嵌入式平台构建RTMP库的开发者&#xff0c;可帮助解决编译链配置与Makefile调整…

作者头像 李华
网站建设 2026/9/9 19:19:44

苏州Android工程师岗位深度拆解:从应用层到系统层面试准备

周六晚上刷到一条苏州虹保世纪科技的 Android 开发工程师岗位推送&#xff0c;JD 写得挺实诚&#xff0c;技术栈列得算清楚&#xff0c;待遇区间也标了范围。我把这家公司近两三年放出来的 Android 相关职位、面试反馈和内部工具链信息翻了个遍&#xff0c;也对照了苏州本地同类…

作者头像 李华
网站建设 2026/9/9 19:19:09

共享储能与多类型负荷需求响应的园区经济运行优化研究

做园区级能量管理项目的时候&#xff0c;我一直有一个很深的感受&#xff1a;很多园区把储能、负荷、光伏当成三个割裂的模块来调度&#xff0c;储能利用率低、负荷侧白白浪费了调节空间&#xff0c;最后算经济账非常难看。这个“含共享储能的园区多类型负荷需求响应经济运行研…

作者头像 李华