Great Expectations 数据质量校验:从一条脏数据到拦截它的 3 类 Expectation
【免费下载链接】great_expectationsAlways know what to expect from your data.项目地址: https://gitcode.com/GitHub_Trending/gr/great_expectations
周一早上,ETL 挂了。原因很简单:上游把user_id从纯数字改成了带前缀的字符串,下游所有CAST全挂了。这类问题,靠人肉看数据根本防不住——字段类型变了没人通知,空值占比悄悄涨了,你只能被动等下游炸。
Great Expectations(简称 GX)就是干这件事的:把“数据应该长什么样”写成可执行的检查规则,数据不达标时让流水线直接停下来。它是开源的 Python 包,面向数据工程师和分析师,核心思路是给数据做单元测试。
Great Expectations 是什么:给数据写的单元测试
一句话定位:GX 是数据质量测试框架,把业务规则翻译成自动执行的校验。
它面向三类人:
- 数据工程师:在入湖、入仓环节加卡点,脏数据别想进来
- 分析师:报表数据对不对,跑一次校验就有答案
- 平台团队:把校验结果做成文档,全组对数据质量有统一认知
和写业务代码的测试一样,GX 里的校验叫Expectation(期望)。区别只在于:被测对象不是函数,而是数据本身。比如“passenger_count必须在 1 到 6 之间”,就是一条 Expectation。规则通过,数据放行;规则失败,立刻拿到失败明细,不用等下游报错再回头查。
运行机制:一个 Batch 如何被校验
GX 的运转围绕四个对象,理解它们的关系,就理解了整个系统。
- Data Context:一切入口,管理数据源、规则集、校验结果这些配置
- Data Source + Batch:描述“数据在哪、取哪一批”。你声明一次连接,GX 用 Batch Request 按批次取数
- Expectation Suite:一组规则的集合,可以按表、按主题组织
- Checkpoint:执行的触发器,把“哪批数据”和“哪套规则”绑在一起
一次校验的完整链路是这样的:
- Batch Request 发出,对应的数据连接器负责取数——Pandas 走内存、SQL 走数据库、Spark 走集群,由执行引擎屏蔽差异
- Validator 拿到这批数据,逐条执行 Expectation Suite 里的规则
- 产出 Validation Result:每条规则通过与否、失败样例、整体通过率
- 触发 Validation Action:更新 Data Docs、发邮件、推 Slack,或者执行任意 Python 逻辑
几个关键设计值得注意:
- 规则和数据源解耦:同一套 Expectation 可以先在测试库跑、再挂到生产表,不用改规则
- 执行引擎统一口径:同一条规则在 Pandas、Spark、SQL 上结果一致,迁移执行环境不用重写逻辑
- 动作可插拔:校验失败后的通知、文档更新都是可配置的动作,不是写死的
这些对象的源码在 核心模块,校验执行逻辑在 校验引擎源码。
如何安装 Great Expectations 并跑通第一条校验
最短路径:装包 → 建上下文 → 连数据 → 加规则 → 跑校验。当前版本 1.22.0,支持 Python 3.10 到 3.13。
先安装:
pip install great_expectations再跑通一条端到端校验,核心代码就这些:
import great_expectations as gx context = gx.get_context() ds = context.data_sources.add_sqlite("demo", sqlite_file_path="taxi.db") asset = ds.add_table_asset(name="rides", table_name="nyc_taxi") batch_definition = asset.add_batch_definition_whole_table("all") suite = context.suites.add(gx.core.expectation_suite.ExpectationSuite(name="rides")) suite.add_expectation(gx.expectations.ExpectColumnValuesToBeBetween( column="passenger_count", min_value=1, max_value=6 )) vd = context.validation_definitions.add( gx.core.validation_definition.ValidationDefinition( name="demo", data=batch_definition, suite=suite ) ) result = context.checkpoints.add( gx.checkpoint.checkpoint.Checkpoint(name="cp", validation_definitions=[vd]) ).run() print(result.describe())result.describe()会打印每条规则的结果和通过率。到这里,第一个校验就跑通了。想换成 Postgres、Snowflake,只需把add_sqlite换成对应的连接器,其余代码不变。数据源连接器一共 40 多个,覆盖主流库和对象存储,实现在 数据源模块。
常用的 Expectation 类型有哪些
GX 内置 60 多种 Expectation,按用途挑最常用的几类:
| 校验类型 | 代表 Expectation | 适用场景 |
|---|---|---|
| 值域 | expect_column_values_to_be_between | 金额、计数不越界 |
| 格式 | expect_column_values_to_match_regex | 邮箱、单号、ID 格式 |
| 空值 | expect_column_values_to_not_be_null | 关键字段不允许空 |
| 唯一性 | expect_column_values_to_be_unique、expect_compound_columns_to_be_unique | 主键、组合键去重 |
| 聚合 | expect_column_mean_to_be_between、expect_table_row_count_to_be_between | 均值、行数波动 |
| 窗口 | expect_column_staleness_to_be_between | 分区数据新鲜度 |
| 分布 | expect_column_values_to_be_in_set | 枚举值白名单 |
两条使用建议:
- 从关键列开始:主键非空、唯一,关键字段值域,先覆盖这三类,能拦下大部分线上事故
- 阈值留余量:行数、均值这类规则设区间而不是固定值,避免正常波动误报
全部规则及参数说明在 官方文档 的 expectation gallery 里,每个类型都有独立示例。
三个典型应用场景
场景一:数据入湖前拦截脏数据。做法:给每张入湖表建一套 Expectation,主键非空唯一、关键字段值域和格式全部写进去,再用 Checkpoint 在入湖任务末尾触发。失败时流水线报错停住,数据不会落到下游。这是性价比最高的用法,因为问题被拦在最便宜的位置。
场景二:监控数据分布漂移。做法:对核心指标列配聚合类规则(均值区间、行数区间),每周跑一次 Checkpoint,结果自动进 Data Docs。某月指标分布明显变化时,报告里一眼能看到哪条规则通过率掉了,不用等报表口径对不上才回查。
场景三:数据迁移后的一致性确认。做法:迁移前后对同一批数据跑聚合校验(行数、总和、均值),源端和目标端各跑一遍做对比;再用expect_compound_columns_to_be_unique确认目标端没有产生重复。迁移是否干净,用数字说话,而不是抽几条看看。
小结
GX 的思路不复杂:规则是 Expectation,执行靠 Checkpoint,结果沉淀在 Data Docs。它不替你做数据架构决策,但能保证“数据不对”这件事第一时间暴露出来。
如果项目里已经有数据管道,建议从一张最常被下游抱怨的表开始试。用过之后有什么心得,欢迎评论区聊聊。
【免费下载链接】great_expectationsAlways know what to expect from your data.项目地址: https://gitcode.com/GitHub_Trending/gr/great_expectations
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考