Verificări ale calității datelor

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

De ce contează porțile de calitate

 

  • Un pipeline fără verificări este un pipeline în care nu poți avea încredere
  • O modificare upstream introduce valori nule sau negative
  • Fiecare dashboard din aval afișează cifre greșite
  • Porțile de calitate prind problemele la sursă

Date proaste care ajung în aval

Construirea pipeline-urilor de date cu 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}}, }, )
Construirea pipeline-urilor de date cu Airflow

Tipuri de verificări pe coloane

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: validează valoarea minimă dintr-o coloană
  • null_check: numără valorile nule (folosește equal_to: 0 pentru a verifica absența lor)
  • distinct_check: numără valorile distincte
  • Verificări disponibile: null_check, unique_check, distinct_check, min, max
  • Comparatori disponibili: equal_to, greater_than, geq_to, less_than, leq_to
Construirea pipeline-urilor de date cu 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", }, }, )
  • Validează proprietăți ale întregului tabel
  • check_statement: orice expresie SQL care returnează adevărat sau fals
Construirea pipeline-urilor de date cu Airflow

Înlănțuirea verificărilor de calitate

@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
  • Încarcă mai întâi, apoi validează
  • Dacă o verificare eșuează, pipeline-ul se oprește
  • Consumatorii din aval nu văd niciodată date greșite
Construirea pipeline-urilor de date cu Airflow

Verificări pe coloane vs. pe tabel

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}},
    },
)
  • Există valori nule?
  • Valorile sunt în intervalul așteptat?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Tabelul conține înregistrări?
  • Totalul se încadrează în limitele așteptate?
Construirea pipeline-urilor de date cu Airflow

Monitorizarea calității la scară cu Astro

$$

  • Astro Observe de la Astronomer ajută la monitorizarea sănătății pipeline-ului
  • Urmărirea SLA cu alerte înainte de termene
  • Genealogia pipeline-ului între DAG-uri și tabele
  • Monitorizarea calității datelor la scară
  • Analiză a cauzelor principale bazată pe AI

Genealogia în Astro Observe

Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...