跳到主要內容

DE Day 28 Airflow 實作:管線上線

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。本地測試推薦三步驟:

  1. 語法檢查:用 python daily_orders.py 跑一次,確認沒有 import 錯誤或語法錯誤。
  2. list DAG:用 docker compose exec scheduler airflow dags list 確認 DAG 被掃描到。
  3. 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 管線綽綽有餘,但放上正式環境時有幾個調校重點:

  1. worker 數量。預設只跑一個 worker,並行能力有限。實務上會啟動 3 到 5 個 worker container。
  2. CeleryExecutor vs KubernetesExecutor。本篇用的是 CeleryExecutor(透過 Redis 分派 Task),適合固定規模。KubernetesExecutor 可以為每個 Task 開一個 Pod,適合資源需求變化大的場景。
  3. 資料庫連線池。當 DAG 很多、Task 都很短時,Postgres 連線會成為瓶頸。可以把 sql_alchemy_pool_size 調大。
  4. 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 設計與除錯模式。

留言

這個網誌中的熱門文章

Day 2 變數與資料型別

Day 2 變數與資料型別 引言 寫程式的過程中,變數與資料型別是處理資料的基礎。變數是存放資料的容器,資料型別則決定這筆資料有哪些特性、可以進行哪些操作。學會定義變數、認識各種資料型別,是學好 Python 的關鍵一步。 這篇文章會帶你了解 Python 中變數的觀念、如何定義變數,以及常見的資料型別,包括整數、浮點數、字串、布林值,還有串列、元組、字典與集合等容器型別。我們也會介紹變數的命名規則與撰寫風格建議,以及如何用 type() 檢查資料型別。 什麼是變數?如何在 Python 中定義變數 變數是在程式執行時用來存放資料的名稱。透過定義變數,我們可以給一筆資料一個名字,並在程式的其他地方用這個名字取用該筆資料。在 Python 中,變數不需要事先宣告型別,因為 Python 是動態型別語言,變數的型別由指定給它的值決定。 定義變數的基本語法 在 Python 中定義變數非常簡單,只要用賦值符號 = 把值指定給變數即可。例如: x = 5 # 定義變數 x,並把整數 5 賦值給它 name = "Alice" # 定義變數 name,並把字串 "Alice" 賦值給它 在這裡,x 是一個變數,被賦予整數 5;name 是另一個變數,被賦予字串 "Alice"。 變數的更新與覆寫 變數的值可以修改,也就是說,我們可以在程式的不同地方給同一個變數新的值。例如: x = 10 # x 最初被賦予 10 x = 15 # x 的值現在被更新為 15 這樣就能依照需求,在程式執行過程中靈活調整變數的值。 Python 的動態型別系統 Python 和某些靜態型別語言不同,定義變數時不需要宣告型別。賦值時,Python 會根據值自動判斷變數的型別。例如: x = 5 # x 是整數 x = 3.14 # x 變成浮點數 x = "Hi" # x 變成字串 同一個變數在程式執行過程中可以存放不同型別的值,這是 Python 的彈性之一。 常見資料型別 在 Python 中,資料型別決定我們可以對變數進行哪些操作...

Day 1 Python 簡介與環境設定

Day 1 Python 簡介與環境設定 引言 在現在的科技環境裡,程式設計已經是一項重要技能。無論你是對資料科學有興趣、想成為開發者,或是想踏入人工智慧(AI)領域,學會寫程式都能明顯提升你的競爭力。在眾多程式語言中,Python 因為語法簡單、功能強大、應用範圍廣泛,成為許多人進入程式世界的第一選擇。這篇文章會帶你認識 Python 的背景與優勢,並一步步教你在不同系統上安裝與設定 Python 開發環境,最後寫出第一支 Python 程式。 為什麼選擇 Python? Python 是一種高階程式語言,由 Guido van Rossum 在 1991 年發布。Python 的設計哲學強調程式碼的可讀性,並用縮排來定義程式區塊,這點和許多使用大括號的語言不同。簡潔的語法讓它成為初學者的理想選擇;就算是經驗豐富的開發者,也能用它完成複雜的專案。 Python 的優勢如下: 簡單易學 :Python 的語法清楚、結構簡潔,初學者很快就能上手。和其他語言相比,學習曲線相對平緩,不需要先弄懂一堆複雜觀念,就能開始寫程式。 應用範圍廣泛 :從資料科學、網頁開發、人工智慧、機器學習、自動化測試到網路爬蟲,Python 都有大量開源函式庫與工具支援,而且在這些領域都扮演關鍵角色。 豐富的函式庫與框架 :Python 的函式庫生態系非常龐大。做資料分析有 NumPy、Pandas;開發網站有 Django、Flask;做深度學習有 TensorFlow、PyTorch。各種需求幾乎都能找到對應的套件,讓開發更有效率。 跨平台支援 :Python 支援 Windows、macOS、Linux 等作業系統,程式通常不需要太多修改就能跨平台執行,讓開發與部署更有彈性。 活躍的社群 :Python 擁有龐大的開發者社群。學習或開發上遇到問題,幾乎都能在社群與論壇(例如 Stack Overflow)找到答案,對初學者來說是很強的後盾,也能減少卡關時的挫折感。 Python 的應用領域 Python 的流行與強大功能,讓許多領域都開始大量使用它。以下是幾個常見的應用方向: 資料科學 :隨著大數據與人工智慧興起,資料科學大量使用 Python。NumPy、Pandas 與 Matplotlib 等工具能處理和分析龐...

Python 從入門到 PyTorch 深度學習:開啟 AI 世界的大門

Python 從入門到 PyTorch 深度學習:開啟 AI 世界的大門 隨著人工智慧(AI)與深度學習(Deep Learning)快速發展,越來越多人對這些技術產生興趣。不論你是想踏入 AI 領域的初學者,還是已經有程式基礎的開發者,學好 Python 與深度學習框架(例如 PyTorch),都能為你打開更多可能。 為什麼選擇 Python? Python 已經是資料科學與人工智慧領域的首選語言。它的語法簡潔、容易上手,而且擁有龐大的生態系與大量開源函式庫。無論是資料處理、資料視覺化,還是建立機器學習與深度學習模型,Python 都能勝任。對想進入 AI 或資料科學領域的人來說,它幾乎是必備工具。 PyTorch 是什麼? PyTorch 是由 Meta(原 Facebook)AI 研究團隊開發的開源深度學習框架,以易用、靈活和動態計算圖著稱,是許多 AI 研究人員與開發者的首選。相較於其他框架,PyTorch 的寫法更貼近原生 Python,對初學者相對友善。無論是簡單的實驗,還是複雜的深度學習模型,PyTorch 都能提供強大的支援。 這個系列能帶給你什麼? 這個系列會從 Python 的基礎開始,帶你一步一步學習,最後能自己用 PyTorch 建立深度學習模型。即使你完全沒有寫過程式,也能跟著文章的節奏累積技能,理解 AI 與深度學習的核心觀念。 本系列涵蓋的主題 Python 基礎:從變數、條件判斷到函式與模組。 資料處理工具:用 NumPy 與 Pandas 有效率地操作資料。 資料視覺化:用 Matplotlib 與 Seaborn 把資料畫成圖表。 深度學習的數學基礎:線性代數、微積分與機率。 PyTorch 入門:理解張量、模型建構與 GPU 加速。 基礎深度學習模型:CNN 與 RNN 的實作應用。 深度學習專案實戰:從資料前處理到模型部署的端到端流程。 誰適合這個系列? 程式初學者 :如果你對 AI 充滿好奇,卻還沒寫過程式,系列的第一部分會帶你快速上手 Python,並幫助你理解深度學習的基本觀念。 資料科學愛好者 :如果你已經熟悉一些資料處理方法,進階部分會教你如何用 PyTorch 建構深度學習模型。 開發者與研究人員 :想更深入了...