Kontrole jakości danych

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Dlaczego bramki jakości są ważne

 

  • Potok bez kontroli to potok, któremu nie można ufać
  • Jedna zmiana upstream wprowadza wartości null lub ujemne
  • Każdy downstream dashboard pokazuje błędne liczby
  • Bramki jakości wychwytują problemy u źródła

Złe dane przepływające downstream

Budowanie potoków danych z 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}}, }, )
Budowanie potoków danych z Airflow

Rodzaje kontroli kolumn

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: sprawdza minimalną wartość w kolumnie
  • null_check: liczy wartości null (użyj equal_to: 0, by wykluczyć nulle)
  • distinct_check: liczy unikalne wartości
  • Dostępne kontrole: null_check, unique_check, distinct_check, min, max
  • Dostępne operatory: equal_to, greater_than, geq_to, less_than, leq_to
Budowanie potoków danych z 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", }, }, )
  • Sprawdza właściwości całej tabeli
  • check_statement: dowolne wyrażenie SQL zwracające prawdę lub fałsz
Budowanie potoków danych z Airflow

Łączenie kontroli jakości w łańcuch

@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
  • Najpierw wczytaj, potem zweryfikuj
  • Gdy kontrola nie przejdzie, potok zatrzymuje się
  • Dalsze etapy nigdy nie zobaczą błędnych danych
Budowanie potoków danych z Airflow

Kontrole kolumn a kontrole 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}},
    },
)
  • Czy są wartości null?
  • Czy wartości mieszczą się w dopuszczalnym zakresie?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Czy tabela ma wiersze?
  • Czy łączna wartość mieści się w oczekiwanych granicach?
Budowanie potoków danych z Airflow

Monitorowanie jakości na dużą skalę z Astro

$$

  • Astro Observe firmy Astronomer pomaga monitorować kondycję potoków
  • Śledzenie SLA z alertami przed upływem terminów
  • Lineage potoków w DAG-ach i tabelach
  • Monitorowanie jakości danych na dużą skalę
  • Analiza przyczyn źródłowych oparta na AI

Lineage w Astro Observe

Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...