Veri kalitesi kontrolleri

Airflow ile Veri İş Hatları Oluşturma

Volker Janz

Senior Developer Advocate at Astronomer

Kalite kapıları neden önemli

 

  • Kontrolsüz bir pipeline, güvenemeyeceğin bir pipeline'dır
  • Tek bir upstream değişiklik null veya negatif değerler sokar
  • Tüm downstream panolar yanlış sayılar gösterir
  • Kalite kapıları sorunları kaynağında yakalar

Kötü verinin aşağı akması

Airflow ile Veri İş Hatları Oluşturma

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}}, }, )
Airflow ile Veri İş Hatları Oluşturma

Sütun kontrol türleri

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: bir sütundaki en küçük değeri doğrular
  • null_check: null değerleri sayar (null olmaması için equal_to: 0 doğrula)
  • distinct_check: farklı değerleri sayar
  • Mevcut kontroller: null_check, unique_check, distinct_check, min, max
  • Mevcut karşılaştırıcılar: equal_to, greater_than, geq_to, less_than, leq_to
Airflow ile Veri İş Hatları Oluşturma

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", }, }, )
  • Tüm tablonun özelliklerini doğrular
  • check_statement: doğru ya da yanlış dönen herhangi bir SQL ifadesi
Airflow ile Veri İş Hatları Oluşturma

Kalite kontrollerini zincirlemek

@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
  • Önce yükle, sonra doğrula
  • Herhangi bir kontrol başarısız olursa pipeline durur
  • Downstream kullanıcıları kötü veriyi asla görmez
Airflow ile Veri İş Hatları Oluşturma

Sütun vs tablo kontrolleri

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}},
    },
)
  • null var mı?
  • Değerler aralıkta mı?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Tabloda satır var mı?
  • Toplam beklenen sınırlar içinde mi?
Airflow ile Veri İş Hatları Oluşturma

Astro ile ölçekte kaliteyi izlemek

$$

  • Astronomer'ın Astro Observe aracı pipeline sağlığını izlemeye yardımcı olur
  • Son tarih öncesi uyarılarla SLA takibi
  • Dag'ler ve tablolar arasında pipeline soy kütüğü
  • Ölçekte veri kalitesi izleme
  • AI tabanlı kök neden analizi

Astro Observe soy kütüğü

Airflow ile Veri İş Hatları Oluşturma

Hadi pratik yapalım!

Airflow ile Veri İş Hatları Oluşturma

Preparing Video For Download...