DE Day 28 Airflow 實作:管線上線
執行需求:需 Docker。昨天把 Airflow 3.x 的觀念走了一遍,今天實際把它跑起來。我們會用 docker-compose 拉起一組 Airflow 3.x(本機示範),裡面包含 Postgres(後端儲存)、Redis(broker)、Airflow Webserver 與 Scheduler,再把我們前幾天的 dbt 專案掛進去當作一支 DAG。讀完之後,你會知道 docker-compose 的結構、DAG 檔案怎麼放、自訂映像檔怎麼建,以及失敗時怎麼重跑。
引言
Airflow 的部署方式有兩條主流路線:自架(用 pip 安裝在機器上)與容器化(用 docker-compose 或 Kubernetes)。自架簡單但環境容易髒、套件衝突難除錯;容器化則把整個 Airflow 生態系(webserver、scheduler、worker、Postgres、Redis)打包進容器,一行 docker compose up 就能把整套系統拉起來。今天我們走容器化路線。
Airflow 3.x 在 2025 年正式發布後,官方推薦用 Docker 部署。本篇範例以 docker-compose 為基礎,把 Airflow 3.0、Postgres 16、Redis 7 一起拉起來,並把我們前幾天的 dbt 專案掛進去當作一支 DAG。整個流程大約 15 分鐘可以跑完第一輪 docker compose up。如果你的本機沒有 Docker,請先到 docker.com 下載 Docker Desktop(或用 OrbStack、Colima 等替代方案)。
今天的目標是:把昨天定義的「每日訂單管線」實際跑起來。完成後你會在 Airflow UI(localhost:8080)看到 DAG、可以手動觸發、可以查看 log、可以重跑失敗的 Task。
目錄結構與設定檔
建立一個工作目錄 de-d28-airflow,裡面會包含:
de-d28-airflow/
├── dags/
│ └── daily_orders.py
├── dbt_project/
│ ├── dbt_project.yml
│ ├── profiles.yml
│ └── models/
├── logs/
├── plugins/
├── docker-compose.yaml
├── Dockerfile
└── requirements.txt
其中 dags/ 是 Airflow 自動掃描的目錄,放進去的 .py 檔會被當作 DAG 載入。logs/ 是 Task log 的輸出位置。plugins/ 是自訂 Operator、Hook、巨集的放置處。
撰寫 Dockerfile
Airflow 3.x 的官方映像檔 apache/airflow:3.0.0-python3.13 已內建大多數常用套件,但 dbt-duckdb 需要另外安裝。建立自訂映像檔:
FROM apache/airflow:3.0.0-python3.13
USER airflow
COPY requirements.txt /requirements.txt
RUN pip install --no-cache-dir -r /requirements.txt
requirements.txt 內容:
dbt-core>=1.10,<2.0
dbt-duckdb
duckdb>=1.3,<2.0
建構映像檔:
docker build -t de-d28-airflow:1.0 .
# 輸出(節錄):
# Sending build context to Docker daemon
# Step 1/4 : FROM apache/airflow:3.0.0-python3.13
# Step 2/4 : USER airflow
# Step 3/4 : COPY requirements.txt /requirements.txt
# Step 4/4 : RUN pip install --no-cache-dir -r /requirements.txt
用 Python 驗證映像檔內的套件
映像檔建好後,可以用 docker exec 進去確認 dbt 與 duckdb 確實安裝成功:
import subprocess
# 檢查容器內的套件版本
result = subprocess.run(
["docker", "exec", "airflow-webserver", "pip", "show", "dbt-core"],
capture_output=True, text=True,
)
print(result.stdout)
# 輸出:
# Name: dbt-core
# Version: 1.10.x
# ...
# 確認 duckdb 從 Python 也能 import
result = subprocess.run(
["docker", "exec", "airflow-webserver", "python", "-c", "import duckdb; print(duckdb.__version__)"],
capture_output=True, text=True,
)
print(result.stdout)
# 輸出:1.4.x
這段 Python 用 subprocess 呼叫 docker exec,幫助你在部署後快速驗證容器內的套件版本是否符合預期。
用 Python 確認 dbt 專案可被 Airflow 容器讀取
把 dbt 專案 mount 進容器後,可以用 Python 從 Airflow 容器的角度檢查它是否能看到 profiles.yml 與模型:
import subprocess
# 列出容器內的 dbt 專案結構
result = subprocess.run(
["docker", "exec", "airflow-scheduler", "ls", "-la",
"/opt/airflow/dbt_project/"],
capture_output=True, text=True,
)
print(result.stdout)
# 輸出(節錄):
# drwxr-xr-x . .
# drwxr-xr-x .. ..
# -rw-r--r-- 1 airflow airflow 248 dbt_project.yml
# -rw-r--r-- 1 airflow airflow 148 profiles.yml
# drwxr-xr-x 4 airflow airflow 128 models
# drwxr-xr-x 2 airflow airflow 64 seeds
# 用 dbt 指令確認連線 OK
result = subprocess.run(
["docker", "exec", "airflow-scheduler", "dbt", "debug",
"--profiles-dir", "/opt/airflow/dbt_project"],
capture_output=True, text=True,
)
print(result.stdout[-500:])
# 輸出(節錄):All checks passed!
這段程式從 host 端呼叫 docker exec,模擬「scheduler 容器能不能成功跑 dbt」。如果 dbt debug 印出 All checks passed,代表 mount 與 profiles 都正確。
docker-compose.yaml 結構
Airflow 3.x 官方提供的 docker-compose.yaml 是個完整範例。我們精簡成五個服務:
x-airflow-common: &airflow-common
build: .
environment: &airflow-common-env
AIRFLOW__CORE__EXECUTOR: CeleryExecutor
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow
AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@postgres/airflow
AIRFLOW__CELERY__BROKER_URL: redis://redis:6379/0
AIRFLOW__CORE__FERNET_KEY: ""
AIRFLOW__CORE__LOAD_EXAMPLES: "false"
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
services:
postgres:
image: postgres:16-alpine
environment:
POSTGRES_USER: airflow
POSTGRES_PASSWORD: airflow
POSTGRES_DB: airflow
healthcheck:
test: ["CMD", "pg_isready", "-U", "airflow"]
interval: 5s
retries: 5
volumes:
- postgres-data:/var/lib/postgresql/data
redis:
image: redis:7-alpine
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 5s
retries: 5
webserver:
<<: *airflow-common
command: webserver
ports:
- "8080:8080"
depends_on:
postgres:
condition: service_healthy
scheduler:
<<: *airflow-common
command: scheduler
worker:
<<: *airflow-common
command: celery worker
volumes:
postgres-data:
這份 compose 有五個服務:postgres(後端)、redis(broker)、webserver(UI)、scheduler(排程器)、worker(執行 Task)。YAML 開頭的 x-airflow-common 是 anchor,把共用設定抽出來讓其他服務繼承,避免重複。
初始化 Airflow 資料庫
第一次啟動前,要先初始化 Postgres 裡的 Airflow metadata:
docker compose up airflow-init
# 這是首次啟動的特殊命令,會跑 migrations 並建立預設帳號
另一個方式是手動跑:
docker compose run --rm airflow-cli airflow db migrate
docker compose run --rm airflow-cli airflow users create \
--username admin \
--password admin \
--firstname Admin \
--lastname User \
--role Admin \
--email admin@example.com
初始化完成後,把整個 Airflow 拉起來:
docker compose up -d
# 等待所有容器變 healthy(約 30 秒)
docker compose ps
瀏覽器開 http://localhost:8080,用 admin/admin 登入,會看到 Airflow 3.x 的 UI。
用 Python 檢查容器健康狀態
Airflow 啟動後,可以用 Python 檢查容器是否都健康:
import subprocess
import json
result = subprocess.run(
["docker", "compose", "ps", "--format", "json"],
capture_output=True, text=True,
)
services = [json.loads(line) for line in result.stdout.strip().split("\n")]
for s in services:
name = s.get("Name", "")
health = s.get("Health", "no-check")
print(f"{name:30} 健康狀態:{health}")
# 輸出(節錄):
# de-d28-airflow-postgres-1 健康狀態:healthy
# de-d28-airflow-redis-1 健康狀態:healthy
# de-d28-airflow-webserver-1 健康狀態:healthy
# de-d28-airflow-scheduler-1 健康狀態:healthy
這段程式讀 docker compose ps 的 JSON 輸出,逐一列出容器健康狀態。如果有容器不是 healthy,先看 log 再決定要不要重啟。
掛入 DAG 與 dbt 專案
建立 dags/daily_orders.py,把昨天定義的概念實作成實際 DAG:
from datetime import datetime, timedelta
from airflow.sdk import dag, task, Asset
orders_asset = Asset("/opt/airflow/dbt_project/data/orders.parquet")
@dag(
dag_id="daily_orders",
schedule="0 2 * * *",
start_date=datetime(2025, 11, 1),
catchup=False,
retries=2,
retry_delay=timedelta(minutes=5),
tags=["orders", "dbt"],
)
def daily_orders():
@task
def download_orders():
"""下載當日訂單 CSV(這裡用 DuckDB 模擬)"""
import duckdb
con = duckdb.connect("/opt/airflow/dbt_project/warehouse.duckdb")
con.execute("""
CREATE OR REPLACE TABLE raw.orders AS
SELECT * FROM read_csv_auto('/opt/airflow/data/orders_*.csv')
""")
return "/opt/airflow/data/orders.parquet"
@task
def run_dbt_models(csv_path: str):
"""跑 dbt seed 與 dbt run"""
import subprocess
result = subprocess.run(
["dbt", "seed", "--profiles-dir", "/opt/airflow/dbt_project"],
capture_output=True, text=True,
)
print(result.stdout)
result = subprocess.run(
["dbt", "run", "--profiles-dir", "/opt/airflow/dbt_project"],
capture_output=True, text=True,
)
print(result.stdout)
if result.returncode != 0:
raise RuntimeError("dbt run 失敗")
return csv_path
@task(outlets=[orders_asset])
def export_warehouse(csv_path: str):
"""把 warehouse 匯出成 Parquet 作為下游產物"""
import duckdb
con = duckdb.connect("/opt/airflow/dbt_project/warehouse.duckdb")
con.execute(f"COPY (SELECT * FROM main.fact_order_line) TO '{csv_path}' (FORMAT PARQUET)")
return csv_path
@task
def notify(path: str):
"""發送通知(這裡只用 print 模擬)"""
print(f"管線完成,產物:{path}")
path = download_orders()
path = run_dbt_models(path)
path = export_warehouse(path)
notify(path)
daily_orders()
這支 DAG 包含四個 Task:下載、跑 dbt、匯出、通知。下游的 DAG 可以透過 orders_asset 觸發(例如「每週儀表板」DAG 等這個產物出來才跑)。
把 dbt 專案也掛進去。在 docker-compose.yaml 的 airflow-common 加上 volumes:
volumes:
- ./dags:/opt/airflow/dags
- ./dbt_project:/opt/airflow/dbt_project
- ./logs:/opt/airflow/logs
重啟 scheduler 讓它重新掃描 DAG:
docker compose restart scheduler
# 等待 30 秒讓 scheduler 載入 DAG
用 Python 觸發與監看 DAG Run
除了 UI 外,可以用 Airflow 的 CLI 或 Python API 觸發 DAG:
import subprocess
# 用 CLI 觸發 DAG
result = subprocess.run(
["docker", "compose", "exec", "-T", "scheduler",
"airflow", "dags", "trigger", "daily_orders"],
capture_output=True, text=True,
)
print(result.stdout)
# 輸出:
# Created <DagRun daily_orders @ 2025-12-05T...: manual__...: success>
# 列出最近 5 次執行
result = subprocess.run(
["docker", "compose", "exec", "-T", "scheduler",
"airflow", "dags", "list-runs", "-d", "daily_orders", "--limit", "5"],
capture_output=True, text=True,
)
print(result.stdout)
# 輸出(節錄):
# run_id state execution_date
# manual__2025-12-05T... success 2025-12-05T02:00:00+00:00
這段 Python 用 docker compose exec 對 scheduler 容器下指令,觸發 DAG 並列出執行紀錄。CLI 比 UI 更快,適合在 CI 或腳本裡自動化觸發。
在 UI 觸發與監看
瀏覽器到 http://localhost:8080,點進 daily_orders DAG,你會看到四個方塊串成的 Graph View。右上角「Trigger DAG」按鈕可以手動觸發一次:
docker compose logs -f worker
# 觀察 worker 執行 Task 的 log
如果某個 Task 失敗,UI 上會顯示紅色,可以點進去看 log,也可以「Clear Task Instance」讓它重跑。「Clear」只會重跑失敗的 Task,下游 Task 會自動接續(這是 Airflow 的標準行為)。
失敗處理與重跑
實務上,DAG 失敗的原因常見有三種:
- 外部資源暫時不可用:API 503、DB timeout。靠
retries=3與retry_delay重試。 - 資料品質問題:CSV 欄位缺值、格式錯誤。靠 Task 內的驗證邏輯提早 raise。
- 商業邏輯錯誤:dbt 模型算錯。靠測試(昨天 Day 26 的內容)提早抓到。
重跑時有兩個常見動作:
- Clear Task Instance:只重跑指定 Task,下游自動接續。適合「修了一個 bug,想驗證後續流程」。
- Clear DAG Run:重跑整個 DAG Run。適合「從頭來過」。
第三個常用動作是「Mark Success」:把失敗的 Task Instance 標記為成功(跳過實際執行),適合「我知道這個問題不重要,先讓整體往下走」。
Volume 掛載與檔案權限
Airflow 容器需要讀寫幾個目錄:dags(讀)、logs(讀寫)、plugins(讀)。如果掛載時權限沒設好,會遇到「Permission denied」錯誤。我們用 Python 驗證 volume 設定:
import subprocess
# 檢查容器能否寫入 logs 目錄
result = subprocess.run(
["docker", "exec", "airflow-scheduler", "bash", "-c",
"echo 'test' > /opt/airflow/logs/test.log && cat /opt/airflow/logs/test.log"],
capture_output=True, text=True,
)
print(result.stdout)
# 輸出:test
# 確認 dags 目錄能被讀取
result = subprocess.run(
["docker", "exec", "airflow-scheduler", "ls",
"/opt/airflow/dags/"],
capture_output=True, text=True,
)
print(result.stdout)
# 輸出:daily_orders.py
# 清理測試檔
subprocess.run(
["docker", "exec", "airflow-scheduler", "rm",
"/opt/airflow/logs/test.log"],
)
這段程式在部署後馬上驗證 volume 是否正確。如果寫不進 logs,scheduler 會立刻 crash;如果讀不到 dags,DAG 就完全不會被載入。
正式環境部署的調校
本機 docker-compose 適合開發,正式環境(VM、Kubernetes)需要調校幾個地方:
- worker 數量:本機只用一個 worker,正式環境會用 3 到 5 個。可以加
worker2、worker3服務(用 docker-compose scale 或 Kubernetes deployment)。 - Postgres 連線池:DAG 很多時,預設 SQLAlchemy pool size 不夠。可以把
sql_alchemy_pool_size從 5 調到 20。 - log 保留天數:預設 log 會無限堆積。設定
log_retention_days=30自動清理。 - 資源限制:用
mem_limit: 2g與cpus: 2限制每個容器,避免 worker 互相搶資源。
另一個進階做法是把 scheduler 與 webserver 拆到不同機器。Airflow 3.x 的設計鼓勵這種拆分:scheduler 純排程、webserver 純 UI、worker 純執行 Task,三者用 metadata 資料庫通訊。
用 Python 監看 worker 資源
正式環境部署後,可以用 Python 定期監看 worker 容器的資源使用:
import subprocess
import json
result = subprocess.run(
["docker", "stats", "--no-stream", "--format", "json"],
capture_output=True, text=True,
)
for line in result.stdout.strip().split("\n"):
s = json.loads(line)
name = s.get("Name", "")
if "airflow" in name:
mem = s.get("MemUsage", "")
cpu = s.get("CPUPerc", "")
print(f"{name:30} CPU: {cpu:8} MEM: {mem}")
# 輸出(節錄):
# de-d28-airflow-webserver-1 CPU: 0.15% MEM: 412.5MiB / 7.7GiB
# de-d28-airflow-scheduler-1 CPU: 0.42% MEM: 380.2MiB / 7.7GiB
# de-d28-airflow-worker-1 CPU: 1.85% MEM: 520.1MiB / 7.7GiB
這段 Python 讀 docker stats 的 JSON,列出 Airflow 容器的 CPU 與記憶體使用。正式環境會把這個輸出推到 Prometheus、Grafana 之類的監控系統。
用 Python 自動重啟失敗容器
正式環境常見的故障情境是「容器健康檢查失敗」或「worker 失聯」。我們可以寫一支簡單的 Python 監控腳本,自動重啟不健康的容器:
import subprocess
import json
def get_unhealthy_containers() -> list:
"""找出所有 health=unhealthy 的容器"""
result = subprocess.run(
["docker", "ps", "--filter", "health=unhealthy",
"--format", "{{.Names}}"],
capture_output=True, text=True,
)
return [name for name in result.stdout.strip().split("\n") if name]
for name in get_unhealthy_containers():
print(f"重啟 {name} ...")
subprocess.run(["docker", "compose", "restart", name.replace("de-d28-airflow-", "")])
print(f"{name} 已重啟")
# 輸出(節錄):
# 重啟 de-d28-airflow-worker-1 ...
# de-d28-airflow-worker-1 已重啟
這段 Python 找到所有不健康的 Airflow 容器,並自動重啟。在正式環境會搭配 cron 或 systemd timer 定期執行(每 5 分鐘一次),確保服務不中斷。
DAG 開發流程:本地測試 → 部署上線
Airflow 的開發流程跟一般 Python 套件不同:因為 DAG 檔案是「被掃描」的,光跑 python daily_orders.py 不會觸發 Airflow。本地測試推薦三步驟:
- 語法檢查:用
python daily_orders.py跑一次,確認沒有 import 錯誤或語法錯誤。 - list DAG:用
docker compose exec scheduler airflow dags list確認 DAG 被掃描到。 - trigger 測試:用
airflow dags trigger daily_orders手動觸發一次,看 Task 是否成功。
這套流程可以用 Python 寫成腳本,每次改完 DAG 自動跑:
import subprocess
# 步驟一:語法檢查
result = subprocess.run(
["python", "dags/daily_orders.py"],
capture_output=True, text=True,
)
assert result.returncode == 0, f"語法錯誤:{result.stderr}"
print("[OK] 語法檢查通過")
# 步驟二:list DAG
result = subprocess.run(
["docker", "compose", "exec", "-T", "scheduler",
"airflow", "dags", "list"],
capture_output=True, text=True,
)
assert "daily_orders" in result.stdout, "DAG 沒被掃描到"
print("[OK] DAG 已載入")
# 步驟三:trigger
result = subprocess.run(
["docker", "compose", "exec", "-T", "scheduler",
"airflow", "dags", "trigger", "daily_orders"],
capture_output=True, text=True,
)
print(f"[OK] 已觸發:{result.stdout.strip()}")
這段 Python 把三步驟自動化,每次改完 DAG 跑一次就能立刻知道有沒有問題。
Airflow 3.x 的 Task SDK 與新架構
Airflow 3.x 的最大架構變革是把 Task 執行邏輯從核心抽出來,變成獨立的 Task SDK。在 2.x 時代,Task 必須在 Airflow worker 容器內執行;3.x 之後,Task 可以跑在任何支援 Python 的環境。
這個改變帶來兩個實際好處:第一,部署更彈性,worker 可以部署到 Kubernetes、AWS ECS、甚至不同作業系統;第二,DAG 檔案只負責定義依賴,執行邏輯由 SDK 處理,DAG 檔案本身變得更輕、更容易除錯。
在 docker-compose 部署裡,這個架構差異體現在 worker 服務的啟動方式:3.x 仍然用 celery worker 啟動,但底層通訊走 Task SDK。對使用者來說,最大的差異是 3.x 推薦用 TaskFlow API(@dag 與 @task 裝飾器),而不是 2.x 的 with DAG 寫法。
另一個 3.x 的關鍵設計是 Edge Labels:DAG 裡的依賴箭頭可以加上文字說明,例如 download >> transform >> load 在 UI 上會顯示「then」「and」等動詞,讓依賴關係更易讀。
Airflow 3.x 的 Migration 注意事項
從 Airflow 2.x 升級到 3.x 時有幾個重點。第一,Task SDK 改寫:所有自訂 Operator、Hook、Executor 都要檢查是否還支援 3.x 的介面。第二,DAG 檔案不要放 top-level code:3.x 的掃描機制對 top-level code 更嚴格,會把呼叫外部副作用的動作視為錯誤。第三,scheduled intervals 寫法:schedule_interval 參數在新版是 schedule,但舊的寫法仍然相容。第四,TaskFlow API 是推薦寫法:如果新寫 DAG,優先用 @dag 與 @task 裝飾器。
實務上,升級前建議先在一個測試環境跑 1 到 2 週,確認所有 Task 都能正常執行,再搬到正式環境。Airflow 3.x 的 changelog 列出了完整的破壞性變更清單,升級前務必讀過。
用 Python 驗證 Docker volume 設定
在部署 Airflow 之前,可以用 Python 模擬 docker-compose 設定,確認 volume 路徑在 host 與容器之間對應正確:
import yaml
with open("docker-compose.yaml") as f:
compose = yaml.safe_load(f)
# 檢查所有服務的 volume mount 設定
for service, config in compose.get("services", {}).items():
if "volumes" in config:
for vol in config["volumes"]:
print(f"{service}: {vol}")
# 輸出(節錄):
# webserver: ./dags:/opt/airflow/dags
# webserver: ./dbt_project:/opt/airflow/dbt_project
# webserver: ./logs:/opt/airflow/logs
# scheduler: ./dags:/opt/airflow/dags
# scheduler: ./dbt_project:/opt/airflow/dbt_project
# scheduler: ./logs:/opt/airflow/logs
# worker: ./dags:/opt/airflow/dags
# worker: ./dbt_project:/opt/airflow/dbt_project
# worker: ./logs:/opt/airflow/logs
這段 Python 讀 docker-compose.yaml,把每個服務的 volume mount 列出來。確認三個 Airflow 服務都有掛 dags、dbt_project、logs 三個目錄,否則 DAG 載入、dbt 執行、log 寫入會出問題。
用 Python 監看 Airflow DAG 狀態
除了 UI,也可以用 Python 直接讀 Airflow 的 metadata 資料庫,監看所有 DAG 的最新狀態:
import subprocess
# 用 docker exec 在 scheduler 容器內跑 airflow CLI
result = subprocess.run(
["docker", "exec", "-T", "airflow-scheduler",
"airflow", "dags", "list", "--output", "json"],
capture_output=True, text=True,
)
import json
dags = json.loads(result.stdout)
for d in dags[:5]:
print(f"{d['dag_id']:30} paused={d.get('is_paused', 'N/A')}")
# 輸出:
# daily_orders paused=False
# example_bash_operator paused=True
這段程式透過 docker exec 從 scheduler 容器取出 DAG 集合(用 JSON 格式),列出每支 DAG 是否被暫停。在正式環境可以加 cron 每小時跑一次,把結果推到監控系統。
另一個監看手段是讀 Airflow 的 metrics 端點:scheduler 與 webserver 預設有 /metrics 端點提供 Prometheus 格式的指標,例如 airflow_dag_run_duration、airflow_task_instance_failed 等。這些指標可以用 Prometheus 抓取後送到 Grafana 視覺化,是正式環境最推薦的監控方案。
綜合以上,Airflow 3.x 在本機部署的整套流程——從 Dockerfile 撰寫、docker-compose 服務編排、Volume 掛載、DAG 觸發、通知與監控——已經比 2.x 友善很多。但「友善」不等於「簡單」,整體而言 Airflow 仍是資料工程裡最重的編排工具之一。
如果你的管線只有一兩條、排程簡單、不需要複雜依賴,Airflow 反而是殺雞用牛刀。明天 Day 29 我們會示範輕量替代方案,讓你在不需要 Airflow 的場景下也能完成編排工作。
最後提醒一點:這套 docker-compose 環境主要是「讓你在本機體驗 Airflow」。正式環境的部署需要更多調校(worker 數量、Postgres 連線池、log 保留天數等),但觀念與本機相同:scheduler 純排程、worker 純執行、webserver 純 UI,三者透過 metadata 資料庫通訊。Day 36 我們會再深入討論部署議題。
至此,Airflow 3.x 的本機部署、DAG 撰寫、Volume 掛載、通知與監控整套流程都走過了。明天 Day 29 我們會從 Airflow 的重量級設計轉向輕量方案,看看在不需要 Airflow 的場景下,怎麼用更簡單的工具達成同樣的編排目標。
通知與監控
在 Airflow 3.x 裡,通知系統也有改進。新的 notifications 模組讓你不用寫 callback 函式,直接在 DAG 上設定失敗通知的頻道:
@dag(
dag_id="daily_orders",
notifications=[
{
"callback": "slack",
"channel": "#data-alerts",
"on_failure": True,
},
{
"callback": "email",
"to": ["team@example.com"],
"on_failure": True,
"on_success": False,
},
],
)
def daily_orders():
...
這比 2.x 時代的 on_failure_callback 簡潔不少。但要注意:3.x 的 notifications 還在演進,部分功能需要特定 Provider 套件,建議先讀官方說明再決定用哪種寫法。
Airflow 內建支援 Slack 與 Email 通知。建立 Slack Connection:
docker compose exec webserver airflow connections add slack_default \
--conn-type slack \
--conn-password "xoxb-YOUR-TOKEN"
然後在 DAG 加上 callback:
from airflow.providers.slack.hooks.slack import SlackHook
def task_failure_alert(context):
"""Task 失敗時發 Slack 通知"""
hook = SlackHook(slack_conn_id="slack_default")
text = f":red_circle: {context['task_instance'].task_id} 失敗於 {context['ds']}"
hook.call("chat.postMessage", json={"channel": "#data-alerts", "text": text})
@dag(
dag_id="daily_orders",
on_failure_callback=task_failure_alert,
...
)
Email 通知類似:把 on_failure_callback 換成寄送 SMTP 的函式即可。Airflow 設定檔裡加上 SMTP 設定。
常見錯誤與踩雷
錯誤一:volume mount 路徑對不起來。DAG 檔案裡寫 /opt/airflow/dags/daily_orders.py,但 compose 裡 mount 到 ./dags:/opt/airflow/dags,這樣才對得起來。常見錯誤是 host 路徑寫成絕對路徑但容器內是相對路徑。
錯誤二:DAG 沒有被載入。檢查三件事:第一,檔案副檔名是 .py 不是 .py.txt;第二,檔案裡沒有 Python 語法錯誤(用 docker compose exec scheduler python /opt/airflow/dags/daily_orders.py 驗證);第三,scheduler 已經掃描過(預設 5 分鐘掃一次,或手動重啟)。
錯誤三:worker 沒安裝套件。pip install 只裝在 webserver 容器內,worker 沒有 dbt 套件。解法是把所有服務都用自訂映像檔啟動,確保每個容器都有相同套件。
錯誤四:忘了設定 FERNET_KEY。Airflow 用 Fernet key 加密連線密碼,預設空字串會導致「Could not create Fernet object」錯誤。可以用 python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())" 產生一個。
錯誤五:Postgres 沒初始化就啟動。第一次啟動會跑 migrations,但如果 Postgres 已經有舊資料,可能會衝突出錯。建議每次重建環境時用 docker compose down -v 把 volumes 一起砍掉重來。
效能與實務提醒
這套 docker-compose 在本機跑 DBT 管線綽綽有餘,但放上正式環境時有幾個調校重點:
- worker 數量。預設只跑一個 worker,並行能力有限。實務上會啟動 3 到 5 個 worker container。
- CeleryExecutor vs KubernetesExecutor。本篇用的是 CeleryExecutor(透過 Redis 分派 Task),適合固定規模。KubernetesExecutor 可以為每個 Task 開一個 Pod,適合資源需求變化大的場景。
- 資料庫連線池。當 DAG 很多、Task 都很短時,Postgres 連線會成為瓶頸。可以把
sql_alchemy_pool_size調大。 - DAG 檔案數量。當 DAG 數量成長到 100 以上,scheduler 載入時間會明顯增加。可以把 DAG 動態生成(用 factory pattern),而不是每個 DAG 一個檔案。
另一個常見提醒:不要在容器內儲存狀態。DAG 的執行結果、log、中間資料都應該寫到外部(S3、Postgres、DuckDB 檔案掛載到 host volume)。容器隨時可以被砍掉重建。
小結
今天我們把 Airflow 3.x 用 docker-compose 拉起來,並掛入一支包含「下載 → 跑 dbt → 匯出 → 通知」四個 Task 的 DAG。整個流程涵蓋了 Dockerfile 自建、docker-compose 結構、Postgres/Redis 依賴、Volume 掛載、UI 操作、失敗重跑、Slack 通知等實際操作面。
這套環境可以延續到後面的 Day 30–35(端到端管線),在那幾篇裡我們會把 Day 18 的政府開放資料、Day 25 的 dbt 專案、Day 27 的 Airflow 概念全部串起來,形成一個可以每天自動跑的資料管線。
結語
Airflow 3.x 的部署已經比 2.x 友善不少,但 docker-compose 加上 Postgres 與 Redis 仍然不是最輕量的方案。如果你的管線只有一兩條、排程簡單、不需要複雜依賴,Airflow 反而是殺雞用牛刀。明天 Day 29 我們會示範輕量替代方案:用 APScheduler、Prefect、Dagster 或純 Python 排程,達成 80% 的效果但只要 20% 的成本。
如果你想保留 Airflow 但簡化部署,也可以考慮官方提供的 Helm chart(Kubernetes)或 Astro(Astronomer 公司維護的託管服務)。
延伸資源
- Apache Airflow 官方文件「Running Airflow in Docker」段落,本篇 compose 檔的官方出處。
- 《Apache Airflow Best Practices》by Marc Lamberti(2023),從單機部署到 Kubernetes 的演進指南。
- Astronomer 部落格「DAG Writing Best Practices」系列,深入討論 DAG 設計與除錯模式。
留言
張貼留言