Kontrola kvality dat

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Proč záleží na quality gates

 

  • Pipeline bez kontrol je pipeline, které nemůžeš věřit
  • Jedna změna na vstupu zavede nuly nebo záporné hodnoty
  • Každý navazující dashboard zobrazí špatná čísla
  • Quality gates zachytí problémy u zdroje

Špatná data šířící se dál

Tvorba datových pipeline s 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}}, }, )
Tvorba datových pipeline s Airflow

Typy kontrol sloupců

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: ověřuje minimální hodnotu ve sloupci
  • null_check: počítá null hodnoty (pomocí equal_to: 0 zajistíš, že žádné nejsou)
  • distinct_check: počítá unikátní hodnoty
  • Dostupné kontroly: null_check, unique_check, distinct_check, min, max
  • Dostupné komparátory: equal_to, greater_than, geq_to, less_than, leq_to
Tvorba datových pipeline s 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", }, }, )
  • Ověřuje vlastnosti celé tabulky
  • check_statement: libovolný SQL výraz vyhodnocený jako true nebo false
Tvorba datových pipeline s Airflow

Řetězení kontrol kvality

@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
  • Nejdřív načti, pak ověř
  • Pokud některá kontrola selže, pipeline se zastaví
  • Navazující systémy nikdy neuvidí špatná data
Tvorba datových pipeline s Airflow

Kontroly sloupců vs. tabulky

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}},
    },
)
  • Jsou ve sloupcích null hodnoty?
  • Jsou hodnoty v očekávaném rozsahu?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Obsahuje tabulka nějaké řádky?
  • Je celkový součet v očekávaných mezích?
Tvorba datových pipeline s Airflow

Monitorování kvality ve velkém s Astro

$$

  • Astro Observe od Astronomer pomáhá monitorovat zdraví pipeline
  • Sledování SLA s upozorněním před vypršením termínu
  • Linie pipeline napříč DAGy a tabulkami
  • Monitorování kvality dat ve velkém měřítku
  • Analýza příčin selhání pomocí AI

Linie v Astro Observe

Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...