Bắt tay viết một file DAG hoàn chỉnh. Bài này đi qua cấu trúc một file DAG kiểu cổ điển với BashOperator và PythonOperator, cách khai báo default_args, schedule, start_date, catchup và toán tử phụ thuộc, rồi giới thiệu TaskFlow API gọn gàng của Airflow 2.
Hầu hết file DAG đều có 4 phần: (1) import, (2) default_args áp chung cho mọi task, (3) khai báo DAG (id, lịch, ngày bắt đầu), và (4) định nghĩa task + phụ thuộc.
from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.bash import BashOperator # Tham số mặc định áp cho mọi task trong DAG default_args = { "owner": "data-team", "retries": 2, "retry_delay": timedelta(minutes=5), } def transform(): # việc nhẹ; việc nặng đẩy sang Spark/SQL print("Đang transform dữ liệu...") return "done" with DAG( dag_id="etl_demo", default_args=default_args, schedule="@daily", # chạy mỗi ngày start_date=datetime(2026, 1, 1), catchup=False, # không chạy bù lịch sử tags=["demo", "etl"], ) as dag: t1 = BashOperator( task_id="extract", bash_command="echo 'Tải dữ liệu nguồn...'", ) t2 = PythonOperator( task_id="transform", python_callable=transform, ) t3 = BashOperator( task_id="load", bash_command="echo 'Nạp vào kho dữ liệu...'", ) # Khai báo phụ thuộc: extract -> transform -> load t1 >> t2 >> t3
Lưu file vào thư mục dags/, Scheduler sẽ tự phát hiện và hiển thị DAG etl_demo trên UI.
Toán tử phụ thuộc >>: t1 >> t2 nghĩa là t2 chạy sau khi t1 thành công. Có thể viết chuỗi t1 >> t2 >> t3, hoặc rẽ nhánh t1 >> [t2, t3] (t2 và t3 chạy song song sau t1). Chiều ngược lại dùng <<.
schedule: nhận nhiều dạng:
| Giá trị | Ý nghĩa |
|---|---|
"@daily" | Mỗi ngày một lần (nửa đêm). Preset khác: @hourly, @weekly, @monthly. |
"0 6 * * *" | Biểu thức cron — 6h sáng hằng ngày. (phút giờ ngày tháng thứ) |
timedelta(hours=3) | Lặp theo khoảng thời gian cố định kể từ start_date. |
None | Không tự chạy — chỉ trigger thủ công hoặc do DAG khác kích hoạt. |
start_date + catchup: start_date là mốc DAG bắt đầu có hiệu lực. Vì Airflow chạy theo khoảng dữ liệu đã kết thúc, DAG Run đầu tiên xuất hiện sau khi khoảng đầu tiên trôi qua. Đặt catchup=False để Airflow bỏ qua các khoảng quá khứ và chỉ chạy khoảng mới nhất.
default_args để mọi task tự thử lại khi lỗi tạm thời (mạng chập chờn, DB bận). Đây là một lý do task cần idempotent: chạy lại không được làm hỏng dữ liệu.TaskFlow API dùng decorator @dag và @task để viết workflow như code Python thông thường. Bạn gọi hàm task như gọi hàm bình thường; Airflow tự suy ra phụ thuộc từ luồng dữ liệu và tự đẩy giá trị trả về qua XCom — không cần khai báo >> thủ công.
from datetime import datetime from airflow.decorators import dag, task @dag( schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False, tags=["demo", "taskflow"], ) def etl_taskflow(): @task() def extract(): return {"rows": 100} # trả về -> đẩy qua XCom @task() def transform(data): return data["rows"] * 2 @task() def load(n): print(f"Nạp {n} bản ghi") # gọi như Python -> Airflow tự suy ra extract >> transform >> load load(transform(extract())) dag_obj = etl_taskflow()
dags/ (mặc định $AIRFLOW_HOME/dags, trong lab Docker là folder dags/ được mount vào container). Scheduler chỉ quét thư mục này.