Controlli di qualità dei dati

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Perché contano i quality gate

 

  • Una pipeline senza controlli è una pipeline di cui non ti puoi fidare
  • Una modifica a monte introduce null o valori negativi
  • Tutte le dashboard a valle mostrano numeri sbagliati
  • I quality gate intercettano i problemi alla fonte

Dati errati che fluiscono a valle

Creare data pipeline con 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}}, }, )
Creare data pipeline con Airflow

Tipi di controlli di colonna

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: valida il valore minimo in una colonna
  • null_check: conta i valori null (usa equal_to: 0 per assenza di null)
  • distinct_check: conta i valori distinti
  • Controlli disponibili: null_check, unique_check, distinct_check, min, max
  • Comparatori disponibili: equal_to, greater_than, geq_to, less_than, leq_to
Creare data pipeline con 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", }, }, )
  • Valida proprietà dell'intera tabella
  • check_statement: qualsiasi espressione SQL che valuti a true o false
Creare data pipeline con Airflow

Collegare i controlli di qualità

@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
  • Prima carica, poi valida
  • Se un controllo fallisce, la pipeline si ferma
  • Chi sta a valle non vede mai dati errati
Creare data pipeline con Airflow

Controlli di colonna vs di tabella

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}},
    },
)
  • Ci sono null?
  • I valori sono nel range?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • La tabella ha righe?
  • Il totale è entro limiti attesi?
Creare data pipeline con Airflow

Monitorare la qualità su larga scala con Astro

$$

  • Astro Observe di Astronomer aiuta a monitorare la salute delle pipeline
  • Tracciamento SLA con avvisi prima delle scadenze
  • Lineage delle pipeline tra DAG e tabelle
  • Monitoraggio qualità dati su larga scala
  • Analisi della causa radice con AI

Lineage in Astro Observe

Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...