Dựng một Airflow chạy được tại máy bằng Docker Compose chính thức của dự án (LocalExecutor + Postgres). Bài này hướng dẫn từng bước cài đặt, đặt một DAG ETL mẫu vào thư mục dags/, và đúc kết best practices khi viết DAG production.
# tạo thư mục lab và vào đó mkdir airflow-lab && cd airflow-lab # tải compose chính thức (chọn đúng version Airflow 2.x) curl -LfO 'https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml'
.env khai báo AIRFLOW_UID (để file do container ghi ra không bị sai quyền):# thư mục cho DAG, log, plugin mkdir -p ./dags ./logs ./plugins ./config # ghi UID hiện tại vào .env (Linux/macOS) echo "AIRFLOW_UID=$(id -u)" > .env
docker compose up airflow-init
docker compose up -d
| Thành phần | URL / Cách dùng |
|---|---|
| Airflow Web UI | localhost:8080 (airflow / airflow) |
| Đặt DAG | Bỏ file .py vào thư mục ./dags — vài chục giây sau DAG hiện trên UI |
| Tắt lab | docker compose down (thêm -v để xoá cả volume DB) |
airflow-init báo cảnh báo RAM/CPU, hãy tăng tài nguyên trong Docker Desktop trước khi up.Một pipeline ETL tối giản: đọc CSV → transform → ghi kết quả (ở đây mô phỏng, in ra số dòng). Trong thực tế bước nặng nên là một SparkSubmitOperator trỏ vào cụm Spark của bạn.
from datetime import datetime, timedelta from airflow.decorators import dag, task default_args = {"owner": "data-team", "retries": 2, "retry_delay": timedelta(minutes=5)} @dag( dag_id="etl_csv_to_parquet", default_args=default_args, schedule="0 6 * * *", # 6h sáng mỗi ngày start_date=datetime(2026, 1, 1), catchup=False, tags=["etl", "lab"], ) def etl_pipeline(): @task() def extract(): import csv path = "/opt/airflow/dags/data/sales.csv" with open(path) as f: rows = list(csv.DictReader(f)) return {"path": path, "count": len(rows)} @task() def transform(meta): # mô phỏng làm sạch: loại dòng qty âm print(f"Transform {meta['count']} dòng từ {meta['path']}") return meta["count"] @task() def load(n): # mô phỏng ghi Parquet ra kho dữ liệu print(f"Ghi {n} dòng ra Parquet (partition theo ngày)") load(transform(extract())) dag = etl_pipeline()
Phiên bản trigger Spark thay cho bước nặng (thay transform/load bằng một task submit Spark):
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator spark_job = SparkSubmitOperator( task_id="spark_transform", application="/opt/airflow/jobs/etl_job.py", conn_id="spark_default", # khai báo trong Admin -> Connections name="etl_csv_to_parquet", )
retries + retry_delay trong default_args để chịu được lỗi tạm thời; cấu hình cảnh báo (email/Slack) khi task fail thật.8080 và xem các DAG ví dụ.dags/etl_daily.py: viết một DAG chạy hằng ngày (schedule="@daily", catchup=False).SparkSubmitOperator gọi pipeline Spark ghi ra Parquet có partition theo ngày; cấu hình spark_default trong Admin → Connections.retries=2 và một FileSensor (mode reschedule) chờ file đầu vào trước khi extract.SparkSubmitOperator vào cụm Spark trong spark-lab, để Airflow đóng vai điều phối còn Spark lo xử lý dữ liệu — đúng kiến trúc thực tế.