Datakvalitetskontroller

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Varför kvalitetsgrindar spelar roll

 

  • En pipeline utan kontroller är en pipeline du inte kan lita på
  • En uppströmsändring kan introducera null-värden eller negativa värden
  • Alla nedströms dashboards visar då felaktiga siffror
  • Kvalitetsgrindar fångar problem vid källan

Felaktig data flödar nedströms

Bygg datapipelines med 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}}, }, )
Bygg datapipelines med Airflow

Typer av kolumnkontroller

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: validerar minimivärdet i en kolumn
  • null_check: räknar null-värden (använd equal_to: 0 för att kräva inga nulls)
  • distinct_check: räknar distinkta värden
  • Tillgängliga kontroller: null_check, unique_check, distinct_check, min, max
  • Tillgängliga jämförelseoperatorer: equal_to, greater_than, geq_to, less_than, leq_to
Bygg datapipelines med 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", }, }, )
  • Validerar egenskaper hos hela tabellen
  • check_statement: valfritt SQL-uttryck som utvärderas till sant eller falskt
Bygg datapipelines med Airflow

Kedja ihop kvalitetskontroller

@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
  • Ladda först, sedan validera
  • Om en kontroll misslyckas stoppas pipelinen
  • Nedströms konsumenter ser aldrig felaktig data
Bygg datapipelines med Airflow

Kolumn- kontra tabellkontroller

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}},
    },
)
  • Finns det null-värden?
  • Ligger värdena inom förväntat intervall?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Innehåller tabellen rader?
  • Ligger totalen inom förväntat intervall?
Bygg datapipelines med Airflow

Övervaka kvalitet i stor skala med Astro

$$

  • Astronomers Astro Observe hjälper till att övervaka pipelinehälsa
  • SLA-spårning med varningar innan deadlines
  • Pipelinehärledning över DAGs och tabeller
  • Datakvalitetsövervakning i stor skala
  • AI-baserad rotorsaksanalys

Astro Observe-härledning

Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...