टेम्पलेटिंग, आइडेमपोटेंसी, और बैकफिलिंग

Airflow के साथ Data Pipelines बनाना

Volker Janz

Senior Developer Advocate at Astronomer

Airflow में Jinja टेम्पलेट्स

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

 

  • {{ ds }} logical date को YYYY-MM-DD फॉर्मेट में रेंडर करता है
  • हर रन की अपनी तारीख होती है, इसलिए वही टास्क हर रन में अलग आउटपुट देता है
  • केवल templatable fields (जैसे bash_command, sql) में काम करता है
Airflow के साथ Data Pipelines बनाना

आइडेमपोटेंसी की समस्या

 

रन 1 (March 31st):

  • March 31st के लिए sales INSERT करें
  • परिणाम: 3 rows

रन 2 (March 31st दोबारा):

  • March 31st के लिए फिर से INSERT करें
  • परिणाम: 6 rows (डुप्लिकेट!)

री-रन पर डुप्लिकेट rows

Airflow के साथ Data Pipelines बनाना

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

$$

  • आप UPSERT या MERGE तरीके भी अपना सकते हैं

Delete-then-insert

Airflow के साथ Data Pipelines बनाना

Backfilling

 

  • historical dates की एक रेंज को दोबारा प्रोसेस करें
  • Airflow हर logical date पर एक Dag रन बनाता है
  • Idempotency के साथ, backfilling सुरक्षित है

Backfill शेड्यूल पर निर्भर है

Airflow के साथ Data Pipelines बनाना

UI में Backfill

  • Dag को Trigger करें, Backfill चुनें, From और To तारीख सेट करें
  • Missing Runs, Missing and Errored Runs, या All Runs में से चुनें
  • Max Active Runs से parallelism सेट करें और execution order बदलें
  • Run Backfill पर प्रक्रिया शुरू होती है और runs बनते हैं

Airflow backfill UI

Airflow के साथ Data Pipelines बनाना

अभ्यास करते हैं!

Airflow के साथ Data Pipelines बनाना

Preparing Video For Download...