Airflow là công cụ điều phối (orchestration) pipeline: định nghĩa bằng Python các bước cần chạy, thứ tự và lịch chạy; scheduler lo chạy đúng giờ, retry khi lỗi, còn UI cho thấy bước nào hỏng. Airflow không tự xử lý dữ liệu nặng — nó ra lệnh cho Spark, dbt hay warehouse làm việc đó.
- DAG (directed acyclic graph) — một workflow: tập các task có hướng, không có vòng lặp, kèm lịch chạy (
schedule). - Task — một đơn vị việc trong DAG. Mỗi lần DAG chạy sinh ra các task instance có trạng thái riêng (
success,failed,up_for_retry...). - Operator — khuôn mẫu của task:
BashOperator,SQLExecuteQueryOperator,SparkSubmitOperator... Sensor là loại operator chỉ dùng để chờ một điều kiện.
python
from datetime import datetime
from airflow.sdk import dag, task
@dag(schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False,
default_args={"retries": 2})
def daily_sales():
@task
def extract(): ...
@task
def load(rows): ...
load(extract()) # extract >> load
daily_sales()So với cron, Airflow có thêm phụ thuộc giữa các bước, retry, lịch sử từng lần chạy, backfill và cảnh báo.
Lưu ý: code ở top-level của file DAG được parse lại định kỳ. Đặt query database hay gọi API ở top-level sẽ làm chậm cả hệ thống; logic phải nằm trong task.