Templating, idempotency และการ backfill

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Jinja templates ใน Airflow

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

 

  • {{ ds }} แสดงผล logical date ในรูปแบบ YYYY-MM-DD
  • แต่ละ run จะได้วันที่ของตัวเอง ทำให้ task เดียวกัน ผลิตผลลัพธ์ต่างกันในแต่ละ run
  • ใช้งานได้เฉพาะใน templatable fields (เช่น bash_command, sql)
การสร้าง Data Pipeline ด้วย Airflow

ปัญหา idempotency

 

Run 1 (31 มีนาคม):

  • INSERT ข้อมูลยอดขายของวันที่ 31 มีนาคม
  • ผลลัพธ์: 3 แถว

Run 2 (รัน 31 มีนาคมอีกครั้ง):

  • INSERT ข้อมูลยอดขายของวันที่ 31 มีนาคมซ้ำ
  • ผลลัพธ์: 6 แถว (ข้อมูลซ้ำ!)

แถวข้อมูลซ้ำเมื่อ re-run

การสร้าง Data Pipeline ด้วย Airflow

รูปแบบ 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

การสร้าง Data Pipeline ด้วย Airflow

การ backfill

 

  • ประมวลผล วันที่ย้อนหลัง ซ้ำในช่วงที่กำหนด
  • Airflow สร้าง DAG run หนึ่งครั้งต่อ logical date
  • เมื่อใช้ร่วมกับ idempotency การ backfill จะ ปลอดภัย

การ backfill ขึ้นอยู่กับ schedule

การสร้าง Data Pipeline ด้วย Airflow

การ backfill ผ่าน UI

  • Trigger DAG แล้วเลือก Backfill กำหนดวันที่ From และ To
  • เลือกได้ว่าจะรัน Missing Runs, Missing and Errored Runs หรือ All Runs
  • กำหนด parallelism ด้วย Max Active Runs และเปลี่ยน ลำดับการ execute ได้
  • กด Run Backfill เพื่อเริ่มกระบวนการและสร้าง run

UI การ backfill ใน Airflow

การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...