Datenqualitätsprüfungen

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Warum Quality Gates wichtig sind

 

  • Eine Pipeline ohne Checks ist eine Pipeline, der du nicht trauen kannst
  • Eine Änderung upstream bringt Nulls oder negative Werte ein
  • Alle nachgelagerten Dashboards zeigen falsche Zahlen
  • Quality Gates fangen Probleme an der Quelle ab

Schlechte Daten fließen nachgelagert

Data-Pipelines mit Airflow aufbauen

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-Pipelines mit Airflow aufbauen

Spalten-Check-Typen

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: validiert den Minimalwert in einer Spalte
  • null_check: zählt Null-Werte (prüfe equal_to: 0 für keine Nulls)
  • distinct_check: zählt verschiedene Werte
  • Verfügbare Checks: null_check, unique_check, distinct_check, min, max
  • Verfügbare Vergleicher: equal_to, greater_than, geq_to, less_than, leq_to
Data-Pipelines mit Airflow aufbauen

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", }, }, )
  • Validiert Eigenschaften der gesamten Tabelle
  • check_statement: beliebiger SQL-Ausdruck, der zu true oder false auswertet
Data-Pipelines mit Airflow aufbauen

Quality Checks verketten

@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
  • Zuerst laden, dann validieren
  • Wenn ein Check fehlschlägt, stoppt die Pipeline
  • Nachgelagerte Nutzer sehen nie schlechte Daten
Data-Pipelines mit Airflow aufbauen

Spalten- vs. Tabellen-Checks

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}},
    },
)
  • Gibt es Nulls?
  • Liegen Werte im Bereich?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Hat die Tabelle Zeilen?
  • Liegt die Summe in erwarteten Grenzen?
Data-Pipelines mit Airflow aufbauen

Qualität in großem Maßstab mit Astro überwachen

$$

  • Astronomers Astro Observe hilft, die Pipeline-Gesundheit zu überwachen
  • SLA-Tracking mit Alerts vor Deadlines
  • Pipeline-Lineage über DAGs und Tabellen hinweg
  • Data-Quality-Monitoring in großem Maßstab
  • KI-gestützte Root-Cause-Analyse

Astro-Observe-Lineage

Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...