Templating, idempotentie en backfilling

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Jinja-templates in Airflow

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

 

  • {{ ds }} rendert de logische datum in YYYY-MM-DD-formaat
  • Elke run heeft een eigen datum, dus dezelfde taak geeft per run ander output
  • Werkt alleen in templatebare velden (zoals bash_command, sql)
Data-pijplijnen bouwen met Airflow

Het idempotentieprobleem

 

Run 1 (31 maart):

  • INSERT sales voor 31 maart
  • Resultaat: 3 rijen

Run 2 (31 maart opnieuw):

  • INSERT sales voor 31 maart nogmaals
  • Resultaat: 6 rijen (duplicaten!)

Dubbele rijen bij opnieuw draaien

Data-pijplijnen bouwen met Airflow

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

$$

  • Je kunt ook UPSERT- of MERGE-aanpakken gebruiken

Eerst verwijderen, dan invoegen

Data-pijplijnen bouwen met Airflow

Backfilling

 

  • Verwerk een reeks historische datums
  • Airflow maakt één Dag-run per logische datum
  • In combinatie met idempotentie is backfilling veilig

Backfill hangt af van het schema

Data-pijplijnen bouwen met Airflow

Backfill in de UI

  • Trigger een Dag, kies Backfill, stel een From- en To-datum in
  • Kies Missing Runs, Missing and Errored Runs of All Runs
  • Stel parallelisme in met Max Active Runs en wijzig de volgorde van uitvoering
  • Met Run Backfill start het proces en worden runs gemaakt

Airflow backfill-UI

Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...