資料品質檢查

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

為何品質閘重要

 

  • 沒有檢查的 pipeline 是你無法信任
  • 上游一個變更就可能引入 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:任一回傳 true/false 的 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
  • 載入,再驗證
  • 任一檢查失敗,pipeline 會停止
  • 下游使用者不會看到壞資料
使用 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 可監控 pipeline 健康度
  • SLA 追蹤:期限前發送警示
  • 跨 DAG 與資料表的血緣關係
  • 大規模資料品質監控
  • 以 AI 進行根因分析

Astro Observe 血緣圖

使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...