Templating, Idempotenz und Backfilling

Data-Pipelines mit Airflow aufbauen

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 das logische Datum im Format YYYY-MM-DD
  • Jeder Lauf hat sein eigenes Datum, daher erzeugt derselbe Task unterschiedliche Ausgaben pro Lauf
  • Funktioniert nur in template-fähigen Feldern (z. B. bash_command, sql)
Data-Pipelines mit Airflow aufbauen

Das Idempotenz-Problem

 

Lauf 1 (31. März):

  • INSERT für Verkäufe am 31. März
  • Ergebnis: 3 Zeilen

Lauf 2 (erneuter 31. März):

  • INSERT für den 31. März erneut
  • Ergebnis: 6 Zeilen (Duplikate!)

Doppelte Zeilen beim erneuten Lauf

Data-Pipelines mit Airflow aufbauen

Das Muster „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, )

$$

  • Alternativ gehen auch UPSERT- oder MERGE-Ansätze

Löschen und dann Einfügen

Data-Pipelines mit Airflow aufbauen

Backfilling

 

  • Einen Bereich historischer Daten neu verarbeiten
  • Airflow erstellt einen DAG-Run pro logischem Datum
  • Zusammen mit Idempotenz ist Backfilling sicher

Backfill hängt vom Zeitplan ab

Data-Pipelines mit Airflow aufbauen

Backfill in der UI

  • Trigger einen DAG, wähle Backfill, setze From- und To-Datum
  • Wähle Missing Runs, Missing and Errored Runs oder All Runs
  • Parallelität mit Max Active Runs festlegen und Ausführungsreihenfolge ändern
  • Mit Run Backfill startet der Prozess und Runs werden erstellt

Airflow-Backfill-UI

Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...