Template, idempotency và backfilling

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Template Jinja trong Airflow

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

 

  • {{ ds }} hiển thị ngày logic theo định dạng YYYY-MM-DD
  • Mỗi lần chạy có ngày riêng, nên cùng một task tạo đầu ra khác nhau cho mỗi lần chạy
  • Chỉ hoạt động trong các trường hỗ trợ template (như bash_command, sql)
Xây dựng Data Pipeline với Airflow

Vấn đề idempotency

 

Chạy 1 (31/3):

  • INSERT dữ liệu sales cho 31/3
  • Kết quả: 3 hàng

Chạy 2 (chạy lại 31/3):

  • INSERT lại dữ liệu sales cho 31/3
  • Kết quả: 6 hàng (bị trùng!)

Hàng bị trùng khi chạy lại

Xây dựng Data Pipeline với Airflow

Mẫu xóa rồi chèn

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

$$

  • Cũng có thể dùng cách UPSERT hoặc MERGE

Xóa rồi chèn

Xây dựng Data Pipeline với Airflow

Backfilling

 

  • Xử lý lại một dải ngày lịch sử
  • Airflow tạo một lần chạy DAG cho mỗi ngày logic
  • Kết hợp với idempotency, backfilling sẽ an toàn

Backfill phụ thuộc vào lịch chạy

Xây dựng Data Pipeline với Airflow

Backfill trong giao diện UI

  • Trigger một DAG, chọn Backfill, đặt ngày FromTo
  • Có thể chọn chạy Missing Runs, Missing and Errored Runs, hoặc All Runs
  • Có thể đặt độ song song bằng Max Active Runs và đổi thứ tự thực thi
  • Nhấn Run Backfill để bắt đầu và tạo các lần chạy

Giao diện backfill của Airflow

Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...