Templating, idempotens och backfilling

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Jinja-mallar i Airflow

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

 

  • {{ ds }} renderar det logiska datumet i formatet YYYY-MM-DD
  • Varje körning får sitt eget datum, så samma task ger olika utdata per körning
  • Fungerar bara i templatingbara fält (t.ex. bash_command, sql)
Bygg datapipelines med Airflow

Idempotensproblemet

 

Körning 1 (31 mars):

  • INSERT försäljning för 31 mars
  • Resultat: 3 rader

Körning 2 (om-körning 31 mars):

  • INSERT försäljning för 31 mars igen
  • Resultat: 6 rader (dubbletter!)

Duplicerade rader vid om-körning

Bygg datapipelines med Airflow

Mönstret 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, )

$$

  • Går även att använda UPSERT- eller MERGE-metoder

Delete-then-insert

Bygg datapipelines med Airflow

Backfilling

 

  • Bearbeta om ett intervall av historiska datum
  • Airflow skapar en DAG-körning per logiskt datum
  • I kombination med idempotens är backfilling säkert

Backfill beror på schemat

Bygg datapipelines med Airflow

Backfill i gränssnittet

  • Trigga en DAG, välj Backfill, ange ett From- och To-datum
  • Välj att trigga Missing Runs, Missing and Errored Runs eller All Runs
  • Ange parallelism med Max Active Runs och ändra körordningen
  • Med Run Backfill startar processen och körningarna skapas

Airflow backfill-gränssnitt

Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...