Kwaliteitscontroles voor data

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Waarom kwaliteitspoorten tellen

 

  • Een pipeline zonder checks is een pipeline die je niet kunt vertrouwen
  • Eén upstream-wijziging zorgt voor nulls of negatieve waarden
  • Elk downstream-dashboard toont verkeerde cijfers
  • Kwaliteitspoorten vangen problemen bij de bron

Slechte data stroomt downstream

Data-pijplijnen bouwen met 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}}, }, )
Data-pijplijnen bouwen met Airflow

Typen kolomchecks

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: valideert de minimumwaarde in een kolom
  • null_check: telt null-waarden (gebruik equal_to: 0 voor geen nulls)
  • distinct_check: telt unieke waarden
  • Beschikbare checks: null_check, unique_check, distinct_check, min, max
  • Beschikbare vergelijkers: equal_to, greater_than, geq_to, less_than, leq_to
Data-pijplijnen bouwen met 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", }, }, )
  • Valideert eigenschappen van de hele tabel
  • check_statement: elke SQL-expressie die true of false oplevert
Data-pijplijnen bouwen met Airflow

Kwaliteitschecks koppelen

@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
  • Eerst laden, daarna valideren
  • Als een check faalt, stopt de pipeline
  • Downstream-verbruikers zien nooit slechte data
Data-pijplijnen bouwen met Airflow

Kolom- vs. tabelchecks

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}},
    },
)
  • Zijn er nulls?
  • Liggen waarden binnen de marge?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Heeft de tabel regels?
  • Ligt het totaal binnen verwachte grenzen?
Data-pijplijnen bouwen met Airflow

Kwaliteit op schaal monitoren met Astro

$$

  • Astronomer's Astro Observe helpt de gezondheid van pipelines te monitoren
  • SLA-tracking met waarschuwingen vóór deadlines
  • Pipeline lineage over Dags en tabellen
  • Datakwaliteitsmonitoring op schaal
  • AI-gedreven root cause analysis

Astro Observe-lineage

Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...