Templating, idempotenza e backfill

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Template Jinja in Airflow

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

 

  • {{ ds }} rende la data logica in formato YYYY-MM-DD
  • Ogni run ha la sua data, quindi lo stesso task produce output diversi per run
  • Funziona solo nei campi templabili (come bash_command, sql)
Creare data pipeline con Airflow

Il problema dell'idempotenza

 

Run 1 (31 marzo):

  • INSERT delle vendite del 31 marzo
  • Risultato: 3 righe

Run 2 (ri-esecuzione 31 marzo):

  • INSERT di nuovo delle vendite del 31 marzo
  • Risultato: 6 righe (duplicati!)

Righe duplicate alla ri-esecuzione

Creare data pipeline con Airflow

Il pattern delete-then-insert

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, )

$$

  • Puoi anche usare approcci UPSERT o MERGE

Elimina e poi inserisci

Creare data pipeline con Airflow

Backfill

 

  • Rielabora un intervallo di date storiche
  • Airflow crea un Dag run per data logica
  • Insieme all'idempotenza, il backfill è sicuro

Il backfill dipende dalla pianificazione

Creare data pipeline con Airflow

Backfill nell'UI

  • Trigger un Dag, seleziona Backfill, imposta date From e To
  • Puoi scegliere Missing Runs, Missing and Errored Runs o All Runs
  • Imposta il parallelismo con Max Active Runs e cambia l'ordine di esecuzione
  • Con Run Backfill il processo parte e vengono creati i run

Interfaccia backfill di Airflow

Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...