Şablonlama, idempotency ve backfill

Airflow ile Veri İş Hatları Oluşturma

Volker Janz

Senior Developer Advocate at Astronomer

Airflow'da Jinja şablonları

@task.bash
def export_data():
    return "cp /data/sales.csv /export/sales_{{ ds }}.csv"

 

  • {{ ds }} mantıksal tarihi YYYY-MM-DD biçiminde üretir
  • Her çalıştırmanın kendi tarihi vardır, yani aynı görev her çalıştırmada farklı çıktı üretir
  • Yalnızca şablonlanabilir alanlarda çalışır (örn. bash_command, sql)
Airflow ile Veri İş Hatları Oluşturma

Idempotency sorunu

 

Çalıştırma 1 (31 Mart):

  • 31 Mart için satışları INSERT et
  • Sonuç: 3 satır

Çalıştırma 2 (31 Mart'ı tekrar çalıştır):

  • 31 Mart satışlarını yeniden INSERT et
  • Sonuç: 6 satır (çift kayıt!)

Yeniden çalıştırmada yinelenen satırlar

Airflow ile Veri İş Hatları Oluşturma

Sil-sonra-ekle deseni

staged_rows = SQLExecuteQueryOperator(
    task_id="get_staged_rows",
    conn_id="duckdb_default",
    sql="SELECT * FROM staging WHERE date = '{{ ds }}'",
)

load = SQLInsertRowsOperator( task_id="load_sales", conn_id="duckdb_default", table_name="sales", columns=["date", "product", "amount"], preoperator="DELETE FROM sales WHERE date = '{{ ds }}';", rows=staged_rows.output, )

$$

  • UPSERT veya MERGE yaklaşımları da kullanılabilir

Önce sil, sonra ekle

Airflow ile Veri İş Hatları Oluşturma

Backfill

 

  • Bir aralıktaki tarihleri yeniden işle
  • Airflow her mantıksal tarih için bir Dag run oluşturur
  • Idempotency ile birlikte backfill güvenlidir

Backfill zamanlamaya bağlıdır

Airflow ile Veri İş Hatları Oluşturma

Arayüzde Backfill

  • Bir Dag'ı tetikle, Backfill seç, From ve To tarihlerini ayarla
  • Missing Runs, Missing and Errored Runs veya All Runs tetiklenebilir
  • Eşzamanlılığı Max Active Runs ile ayarlayabilir, çalışma sırasını değiştirebilirsin
  • Run Backfill ile süreç başlar ve çalıştırmalar oluşturulur

Airflow backfill arayüzü

Airflow ile Veri İş Hatları Oluşturma

Hadi pratik yapalım!

Airflow ile Veri İş Hatları Oluşturma

Preparing Video For Download...