数据质量检查

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

质量闸门为何重要

 

  • 没有检查的流水线是您无法信任
  • 上游一次变更就可能引入null负值
  • 下游看板都会出现错误数字
  • 质量闸门在源头捕获问题

劣质数据向下游流动

使用 Airflow 构建数据流水线

SQLColumnCheckOperator

from airflow.providers.common.sql.operators.sql import
    (SQLExecuteQueryOperator, SQLColumnCheckOperator)

check_columns = SQLColumnCheckOperator( task_id="check_columns", conn_id="duckdb_analytics", table="daily_summary", column_mapping={ "total_revenue": {"min": {"greater_than": 0}}, "total_orders": {"null_check": {"equal_to": 0}}, }, )
使用 Airflow 构建数据流水线

列级检查类型

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min:校验列的最小值
  • null_check:统计 null 数量(断言 equal_to: 0 表示无 null)
  • distinct_check:统计去重后的数量
  • 可用检查:null_checkunique_checkdistinct_checkminmax
  • 可用比较器:equal_togreater_thangeq_toless_thanleq_to
使用 Airflow 构建数据流水线

SQLTableCheckOperator

from airflow.providers.common.sql.operators.sql import SQLTableCheckOperator

check_table = SQLTableCheckOperator( task_id="check_table", conn_id="duckdb_analytics", table="daily_summary", checks={ "row_count_check": { "check_statement": "COUNT(*) > 0", }, }, )
  • 校验整张表的属性
  • check_statement:任意返回真或假的 SQL 表达式
使用 Airflow 构建数据流水线

串联质量检查

@dag(schedule="@daily", template_searchpath="/path/to/sql")
def sales_pipeline():

    load = SQLExecuteQueryOperator(
        task_id="load_daily_sales",
        conn_id="duckdb_analytics",
        sql="aggregate_sales.sql",
    )

    check_columns = SQLColumnCheckOperator(...)
    check_table = SQLTableCheckOperator(...)

    load >> check_columns >> check_table
  • 加载,再校验
  • 任一检查失败,流水线即停止
  • 下游使用方永不见到坏数据
使用 Airflow 构建数据流水线

列级 vs 表级检查

SQLColumnCheckOperator

check_columns = SQLColumnCheckOperator(
    task_id="check_columns",
    conn_id="duckdb_analytics",
    table="daily_summary",
    column_mapping={
        "total_revenue": {"min": {"greater_than": 0}},
        "total_orders": {"null_check": {"equal_to": 0}},
    },
)
  • 是否存在null
  • 数值是否在范围内

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • 表里有行吗?
  • 汇总是否在预期范围内?
使用 Airflow 构建数据流水线

用 Astro 实现大规模质量监控

$$

  • Astronomer 的 Astro Observe 可监控流水线健康
  • SLA 跟踪:截止前预警
  • 跨 Dag 与表的血缘关系
  • 数据质量监控可规模化运行
  • 基于 AI 的根因分析

Astro Observe 血缘

使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...