Šablonování, idempotence a backfilling

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Jinja šablony v Airflow

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

 

  • {{ ds }} vykreslí logické datum ve formátu YYYY-MM-DD
  • Každý běh má vlastní datum, takže stejný task produkuje různý výstup
  • Funguje pouze v šablonovatelných polích (např. bash_command, sql)
Tvorba datových pipeline s Airflow

Problém idempotence

 

Běh 1 (31. března):

  • INSERT prodejů za 31. března
  • Výsledek: 3 řádky

Běh 2 (opakování 31. března):

  • INSERT prodejů za 31. března znovu
  • Výsledek: 6 řádků (duplikáty!)

Duplicitní řádky při opakovaném běhu

Tvorba datových pipeline s Airflow

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

$$

  • Lze použít i přístupy UPSERT nebo MERGE

Smazání a vložení

Tvorba datových pipeline s Airflow

Backfilling

 

  • Znovu zpracuje rozsah historických dat
  • Airflow vytvoří jeden běh DAGu na každé logické datum
  • V kombinaci s idempotencí je backfilling bezpečný

Backfill závisí na harmonogramu

Tvorba datových pipeline s Airflow

Backfill v uživatelském rozhraní

  • Spusť DAG, vyber Backfill, nastav datum Od a Do
  • Lze zvolit Missing Runs, Missing and Errored Runs nebo All Runs
  • Paralelismus nastavíš přes Max Active Runs, pořadí spouštění lze také změnit
  • Po kliknutí na Run Backfill se proces spustí a běhy se vytvoří

Uživatelské rozhraní backfillu v Airflow

Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...