Bài 2

Viết DAG đầu tiên

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.

Cấu trúc một file DAG
Bốn phần quen thuộc

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.

Cấu trúc file DAG Scroll / zoom · Mở draw.io ↗
DAG kiểu cổ điển
BashOperator + PythonOperator
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.

Giải thích các thành phần
Phụ thuộc, lịch chạy, ngày bắt đầu

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.
NoneKhô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.

retries & retry_delay: đặt trong 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 (Airflow 2)
Viết DAG như viết hàm Python

TaskFlow API dùng decorator @dag@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()
Đặt file ở đâu: mọi file DAG phải nằm trong thư mục 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.
← Bài trước
Tổng quan & Kiến trúc