DE Day 27 Airflow 概念與 DAG 設計
執行需求:CPU 可跑。今天是「編排」區塊的開篇章。我們不裝 Airflow、不跑任何容器,只用「看圖說故事」的方式把 Airflow 3.x 的核心觀念走一遍:DAG 是什麼、Task 之間怎麼串、Operator 與 Sensor 怎麼選、TaskFlow API 怎麼用、Airflow 3.x 新引入的 Asset 概念是什麼。讀完之後,你會對「編排」這個詞有具體的畫面,知道 Day 28 裝 Airflow 時在裝什麼。
引言
Day 26 結束時,我們有一個能跑的 dbt 專案:模型有版本控制、測試有 CI、說明有 dbt docs。但有一個問題還沒解決——誰來每天幫我們跑 dbt run?誰來抓 GitHub 上新的來源檔案、轉換、落地、發通知?靠工程師手動執行不可能,靠 crontab 又會遇到失敗沒人通知、依賴關係不明確、重跑邏輯複雜等問題。Airflow 就是為了把「定時跑資料管線」這件事系統化而誕生的工具。
Airflow 在 2014 年由 Airbnb 開源,2020 年成為 Apache 基金會的頂級專案,是目前業界最普及的編排工具。它在 2025 年 4 月發布了 Airflow 3.x 正式版,帶來不少重要的設計改動:Task SDK(把工作邏輯與排程核心解耦)、Asset 概念(讓 DAG 之間用「資料產物」連接)、Edge Labels(讓依賴關係自帶說明文字)、DAG Versioning(DAG 檔案可版本化)。今天這篇會把這些觀念都點到。
這篇文章不示範任何完整範例的執行(那是 Day 28 的工作),但會用 Python 程式片段展示「DAG 的程式碼長什麼樣」,讓你在裝 Airflow 之前先有心理準備。今天所有程式碼都可以在普通 Python 直譯器跑,不依賴 Airflow 本體。
什麼是 DAG
DAG 是 Directed Acyclic Graph(有向無環圖)。在 Airflow 的語境裡,DAG 是一份 Python 檔案,描述「一組任務與它們之間的執行順序」。每一個 DAG 對應一支資料管線,例如「每日訂單管線」、「每週備份管線」、「每月報表管線」。
幾個關鍵詞要先釐清:
- Task(任務):DAG 裡的最小執行單位,例如「下載 CSV」、「跑 dbt 模型」、「寄送電子郵件」。一個 DAG 通常有 5 到 50 個 Task。
- Operator:Task 的範本。PythonOperator 用來執行 Python 函式、BashOperator 用來跑 shell 指令、PostgresOperator 用來執行 SQL。Operator 是 Task 的「型別」,Task 是 Operator 的「實例」。
- Sensor:特殊的 Operator,專門用來「等待某個條件成立」。例如等檔案出現、等另一個 DAG 完成、等外部 API 釋出新資料。
- Schedule Interval(排程間隔):DAG 多久跑一次。可以是 cron 表達式("0 2 * * *" 每天凌晨兩點)、preset(@daily、@hourly)、或 None(手動觸發)。
- Task Instance(任務實例):某個 Task 在某個具體執行週期裡的執行紀錄。DAG 每天跑一次就會產生一組 Task Instance。
DAG 的兩個核心限制:「有向」是因為任務之間有明確的上下游;「無環」是因為不允許 A 等 B 完成、B 又等 A 完成,否則會無限迴圈。
用 Python 模擬 DAG 的執行順序
在還沒裝 Airflow 之前,我們用 Python 的 dataclass 與簡單的拓樸排序來模擬 DAG 的概念:
from dataclasses import dataclass, field
from typing import List, Dict
@dataclass
class Task:
task_id: str
depends_on: List[str] = field(default_factory=list)
# 模擬一份 DAG 的依賴圖
tasks = {
"download": Task(task_id="download"),
"transform": Task(task_id="transform", depends_on=["download"]),
"load": Task(task_id="load", depends_on=["transform"]),
"test": Task(task_id="test", depends_on=["load"]),
"notify": Task(task_id="notify", depends_on=["test"]),
}
def topological_order(tasks: Dict[str, Task]) -> List[str]:
"""模擬 Airflow 排程器決定 Task 執行順序"""
visited = set()
order = []
def visit(task_id: str):
if task_id in visited:
return
visited.add(task_id)
for dep in tasks[task_id].depends_on:
visit(dep)
order.append(task_id)
for tid in tasks:
visit(tid)
return order
print("執行順序:", topological_order(tasks))
# 輸出:執行順序: ['download', 'transform', 'load', 'test', 'notify']
這段程式把 Airflow 排程器的核心邏輯——拓樸排序——用 20 行 Python 重現。真實的 Airflow 排程器在做的事跟這差不多,只是會把順序記到資料庫、發訊號給 worker、處理失敗與重試。
第一支 DAG:手繪概念到程式碼
假設我們有一條簡單的每日管線:下載訂單 CSV → 跑 dbt 模型 → 寄送通知。用 Airflow 的程式碼寫起來會像這樣:
from datetime import datetime, timedelta
from airflow.sdk import dag, task
@dag(
dag_id="daily_orders",
schedule="0 2 * * *",
start_date=datetime(2025, 11, 1),
catchup=False,
tags=["orders", "daily"],
)
def daily_orders():
@task
def download_orders():
print("下載訂單 CSV ...")
return "/data/orders_20251205.csv"
@task
def run_dbt(csv_path: str):
print(f"用 {csv_path} 跑 dbt ...")
return "dbt run 完成"
@task
def send_notification(dbt_result: str):
print(f"寄信通知:{dbt_result}")
csv_path = download_orders()
dbt_result = run_dbt(csv_path)
send_notification(dbt_result)
daily_orders()
這段程式碼用的是 Airflow 3.x 引入的 TaskFlow API(基於 Python 裝飾器)。幾個關鍵改變:
- @dag 與 @task 裝飾器:3.x 推薦的寫法,比 2.x 的
with DAG(...):寫法簡潔。 - 函式之間用 XCom 自動傳值:
download_orders()回傳的字串會自動被run_dbt(csv_path)接收,Airflow 透過 XCom(cross-communication)機制在底層處理。 - schedule 接受 cron 表達式:在 3.x 裡
schedule="0 2 * * *"等同於schedule_interval="0 2 * * *"。 - catchup=False:不補跑歷史。如果 DAG 從 11 月才開始部署,不會自動把 11 月以前沒跑的排程補上。
Airflow 3.x 把「DAG 定義」與「Task 執行」拆開了。DAG 定義仍然在排程器(scheduler)裡跑,但 Task 的執行邏輯現在透過 Task SDK 在 worker 跑。這讓部署更靈活——可以把 worker 部署到 Kubernetes 或不同機器。
Operator 與 Sensor 的選擇
Airflow 內建 50 種以上的 Operator,常用的有:
- @task:執行 Python 函式(TaskFlow API),3.x 推薦的寫法。
- BashOperator:執行 shell 指令。例如
bash_command="dbt run --profiles-dir ."。 - PythonOperator:2.x 風格的 Python 任務,3.x 仍然支援但推薦改用 @task。
- EmptyOperator:佔位 Task,用於視覺化分組,不做事。
Sensor 方面:
- FileSensor:等檔案出現。例如
filepath="/data/orders.csv"。 - S3KeySensor:等 S3(或 MinIO)上的物件出現。
- ExternalTaskSensor:等另一個 DAG 的某個 Task 完成。
- HttpSensor:定期打 HTTP 端點直到回傳預期狀態。
Sensor 的設計哲學是「polling」(輪詢),它會定期檢查條件是否成立,預設每 60 秒一次,可以調短。Sensor 適合「等待外部事件」,例如等合作方上傳檔案、等 API 釋出資料。但要注意:Sensor 太多會佔用 worker,建議把長時間等待的 Sensor 集中到少數 DAG。
用 dataclass 模擬 Task 與 Sensor 的概念
在實際裝 Airflow 之前,我們先用 Python 類別釐清 Task、Operator、Sensor 的差異:
from dataclasses import dataclass
from abc import ABC, abstractmethod
class Operator(ABC):
"""所有 Operator 的共同祖先"""
@abstractmethod
def execute(self, context): ...
@dataclass
class PythonTask(Operator):
"""@task 裝飾器的概念"""
func: callable
def execute(self, context):
return self.func(context)
@dataclass
class BashTask(Operator):
"""BashOperator 的概念"""
command: str
def execute(self, context):
import subprocess
return subprocess.run(self.command, shell=True, check=True)
@dataclass
class FileSensor(Operator):
"""FileSensor 的概念:等檔案出現"""
filepath: str
poke_interval: int = 60
def execute(self, context):
import os, time
while not os.path.exists(self.filepath):
print(f"等 {self.filepath} ...")
time.sleep(self.poke_interval)
return self.filepath
print("Operator 與 Sensor 的差別在 execute() 的實作")
print("Task 是 Operator 的『包裝』,DAG 是 Task 的『排序』")
這段程式碼展示了 Airflow 內部如何用「物件導向」組織 Operator 與 Sensor。實務上你不會自己寫這些,但讀懂概念對除錯有幫助。
XCom 與參數傳遞
Airflow 用 XCom(cross-communication)在 Task 之間傳遞小資料。預設透過 metadata 資料庫存,3.x 起可以用 Task SDK 直接傳:
from airflow.sdk import dag, task
from datetime import datetime
@dag(
dag_id="xcom_demo",
start_date=datetime(2025, 11, 1),
schedule=None,
)
def xcom_demo():
@task
def producer() -> dict:
"""產生資料:dict、list、字串都可以傳"""
return {
"file_path": "/data/orders_20251205.csv",
"row_count": 12345,
"min_date": "2025-11-01",
}
@task
def consumer(payload: dict) -> None:
"""接收資料"""
print(f"收到檔案:{payload['file_path']}")
print(f"列數:{payload['row_count']}")
data = producer()
consumer(data)
xcom_demo()
XCom 設計上有三條限制:第一,只傳小資料(小於 1 MB),大資料應該寫到 S3 或 DuckDB 後只傳檔名;第二,傳的東西要可序列化(JSON、str、int),不要傳 datetime 物件(會序列化失敗);第三,XCom 預設在 metadata 資料庫裡,大規模使用會拖慢資料庫,需要調整 xcom_backend 設定。
用 Python 模擬 XCom 的底層
XCom 的精神是「Task 之間傳遞小資料」。我們用 dict 模擬:
from typing import Any, Dict, List
class XComStore:
"""模擬 Airflow XCom 的簡單實作"""
def __init__(self):
self.store: Dict[str, Any] = {}
def push(self, key: str, value: Any):
"""producer 把值推入"""
self.store[key] = value
return key
def pull(self, key: str) -> Any:
"""consumer 拉取值"""
return self.store.get(key)
xcom = XComStore()
xcom.push("download_result", {"path": "/data/orders.csv", "rows": 12345})
result = xcom.pull("download_result")
print(f"consumer 收到:{result}")
# 輸出:consumer 收到:{'path': '/data/orders.csv', 'rows': 12345}
這段程式把 XCom 的行為用 Python 重現。真實的 Airflow 會把資料存到 Postgres,並加上 key、task_id、execution_date 等 metadata。
Airflow 3.x 的 Asset 概念
Airflow 3.x 引入了一個新概念叫 Asset(早期叫做 Dataset)。它的想法是:「DAG 之間的依賴不應該用時間排程表,而應該用『資料產物』相連」。
範例:假設我們有「下載訂單」DAG 產出 orders.csv,與「跑分析」DAG 消費 orders.csv。傳統寫法是讓「跑分析」DAG 用 cron 排固定時間,希望「下載訂單」DAG 已經跑完。Asset 寫法則是:
from airflow.sdk import Asset, dag, task
orders_asset = Asset("file:///data/orders.csv")
@dag(schedule=None, tags=["ingest"])
def download_orders():
@task(outlets=[orders_asset])
def download():
print("下載並產出 orders.csv ...")
download_orders()
@dag(schedule=[orders_asset], tags=["analytics"])
def analyze_orders():
@task
def analyze():
print("讀取 orders.csv 跑分析 ...")
analyze_orders()
這段程式碼的重點是 schedule=[orders_asset]。它告訴 Airflow:這個 DAG 不是按時間排程,而是「等 orders_asset 被產出時自動觸發」。當 download_orders 的 @task(outlets=[orders_asset]) 完成後,Airflow 會自動觸發 analyze_orders。
Asset 把「資料管線的依賴」從「時間」拉回到「資料本身」,是 Airflow 3.x 最核心的設計變革之一。實務上,當你有 10 條以上 DAG 時,Asset 能顯著降低「A 還沒跑完,B 就開始跑」的競爭條件。
用 Python 模擬 Asset 觸發邏輯
Asset 的精神是「資料產出後自動通知下游」。我們用簡單的 Python 模擬這個機制:
from dataclasses import dataclass, field
from typing import List, Dict, Callable
@dataclass
class Asset:
uri: str
produced: bool = False
@dataclass
class AssetTask:
name: str
produces: List[Asset] = field(default_factory=list)
consumes: List[Asset] = field(default_factory=list)
assets: Dict[str, Asset] = {
"orders": Asset("file:///data/orders.csv"),
"report": Asset("file:///data/report.parquet"),
}
tasks = [
AssetTask("download_orders", produces=[assets["orders"]]),
AssetTask("analyze_orders", consumes=[assets["orders"]], produces=[assets["report"]]),
AssetTask("send_report", consumes=[assets["report"]]),
]
def run_pipeline(tasks):
"""模擬 Asset 驅動的 DAG 執行"""
pending = list(tasks)
while pending:
progress = False
for t in pending:
if all(a.produced for a in t.consumes):
print(f"執行 {t.name}")
for a in t.produces:
a.produced = True
pending.remove(t)
progress = True
break
if not progress:
print("卡住:某個 Task 的依賴未滿足")
break
run_pipeline(tasks)
# 輸出:
# 執行 download_orders
# 執行 analyze_orders
# 執行 send_report
這段程式展示「Asset 驅動」的執行邏輯:每個 Task 檢查它的 consumes 是不是都已產出,是的話就執行並標記 produces 為已產出。Airflow 內部的 TaskInstance 排程邏輯跟這個非常類似。
連線與 Variable 的概念示範
Airflow 把外部資源的設定集中在 Connection 與 Variable。這個設計讓 DAG 可以「同一支程式、跨環境執行」。我們用 Python 模擬這個概念:
from dataclasses import dataclass
from typing import Dict, Any
@dataclass
class Connection:
conn_id: str
conn_type: str # postgres / s3 / http ...
host: str
port: int
login: str
password: str
schema: str = "public"
# 模擬 dev 與 prod 兩組設定
connections = {
"dev_db": Connection(
conn_id="dev_db", conn_type="postgres",
host="localhost", port=5432,
login="dev", password="dev_pw", schema="dev",
),
"prod_db": Connection(
conn_id="prod_db", conn_type="postgres",
host="prod.example.com", port=5432,
login="etl", password="prod_pw", schema="public",
),
}
variables = {
"batch_size": 10000,
"alert_channel": "#data-alerts",
}
# 在 DAG 裡用 conn_id 與 Variable.get 取得設定
def get_connection(env: str) -> Connection:
return connections[f"{env}_db"]
def get_variable(key: str) -> Any:
return variables[key]
print(f"dev 環境 DB:{get_connection('dev').host}")
# 輸出:dev 環境 DB:localhost
print(f"目前 batch_size:{get_variable('batch_size')}")
# 輸出:目前 batch_size:10000
這段程式把 Airflow Connection 與 Variable 的概念用 dataclass 重現。真實的 Airflow 會把這些資訊存在 metadata 資料庫,DAG 透過 conn_id 與 Variable.get("key") 動態取得。
用 Python 模擬 XCom 的底層
XCom 的精神是「Task 之間傳遞小資料」。我們用 dict 模擬:
from typing import Any, Dict, List
class XComStore:
"""模擬 Airflow XCom 的簡單實作"""
def __init__(self):
self.store: Dict[str, Any] = {}
def push(self, key: str, value: Any):
"""producer 把值推入"""
self.store[key] = value
return key
def pull(self, key: str) -> Any:
"""consumer 拉取值"""
return self.store.get(key)
xcom = XComStore()
xcom.push("download_result", {"path": "/data/orders.csv", "rows": 12345})
result = xcom.pull("download_result")
print(f"consumer 收到:{result}")
# 輸出:consumer 收到:{'path': '/data/orders.csv', 'rows': 12345}
這段程式把 XCom 的行為用 Python 重現。真實的 Airflow 會把資料存到 Postgres,並加上 key、task_id、execution_date 等 metadata。
XCom 在 3.x 有兩個重要的設定要記住:
- xcom_backend:預設是 metadata 資料庫,量大會拖慢;可以改用 S3 或 Redis。
- enable_xcom_pickling:預設關閉,要傳複雜物件(如自訂類別)需要打開,但會有安全風險。
用 Python 檢查 DAG 是否有環狀依賴
Airflow 不允許 DAG 出現環狀依賴。可以用 Python 簡單的 DFS 檢查:
from typing import Dict, List
def has_cycle(graph: Dict[str, List[str]]) -> bool:
"""檢查 DAG 是否出現環狀依賴"""
WHITE, GRAY, BLACK = 0, 1, 2
color = {n: WHITE for n in graph}
def dfs(node):
color[node] = GRAY
for nxt in graph.get(node, []):
if color[nxt] == GRAY:
return True # 發現環
if color[nxt] == WHITE and dfs(nxt):
return True
color[node] = BLACK
return False
for n in graph:
if color[n] == WHITE and dfs(n):
return True
return False
# 模擬 DAG 依賴圖
dag = {
"download": [],
"transform": ["download"],
"load": ["transform"],
"notify": ["load"],
}
print(f"無環:{not has_cycle(dag)}")
# 輸出:無環:True
# 加一個環狀依賴測試
bad_dag = {
"A": ["B"],
"B": ["C"],
"C": ["A"],
}
print(f"無環:{not has_cycle(bad_dag)}")
# 輸出:無環:False
這段 Python 用 DFS 與「灰黑」標記偵測環。Airflow 在排程器載入 DAG 時也會跑類似的檢查,遇到環狀依賴會拒絕載入並印出錯誤訊息。
DAG 設計原則
DAG 設計的另一個容易被忽略的細節是「Task 名稱的命名規則」。Airflow 的 UI 會把 Task 名稱當作識別依據,建議用「動詞_名詞」的格式(例如 download_orders、run_dbt_models、send_notification),不要用模糊的 step1、step2、task_final 這類命名。當 DAG 失敗時,工程師第一眼會看 Task 名稱判斷「是哪裡出問題」,清晰的名稱能省下大量除錯時間。
除了前面列的五條原則,這裡再補充三條進階守則。第一,用 TaskGroup 組織 Task。當 DAG 超過 10 個 Task 時,可以用 with TaskGroup(group_id="etl") as tg: 把相關 Task 包成一組,UI 上會自動變成可收合的子圖。第二,設定合理的 timeout。每個 Task 可以設 execution_timeout=timedelta(minutes=30),避免某個 Task 死結佔住 worker。第三,避免在 DAG 裡寫 print。Airflow 有自己的 log 系統,用 context["ti"].log.info("訊息") 印 log,會自動出現在 UI 上,比 print 更結構化。
好的 DAG 設計有幾個原則,新手一定要記住:
- Task 顆粒度適中。太細(每個 Task 只做一件事)會讓 DAG 變複雜,執行 overhead 也大;太粗(一個 Task 做所有事)又失去編排的意義。一個 Task 的執行時間最好在 5 秒到 30 分鐘之間。
- Task 之間沒有環狀依賴。A 等 B 完成、B 等 A 完成是不可能的。DAG 必須是無環圖。
- 避免在 DAG 檔案裡寫太多商業邏輯。DAG 檔案只負責「定義任務與依賴」,商業邏輯放在獨立的 Python 模組或 dbt 模型。
- 冪等(idempotent)。同一個 Task 跑兩次應該跟跑一次結果一樣。這樣重跑與補跑才安全。
- Task 失敗時要明確處理。
retries=3與retry_delay=timedelta(minutes=5)是常見設定。失敗超過次數後用on_failure_callback發通知。
一個常見的設計錯誤是把所有事情都塞進一個 DAG。每天的訂單管線、每週的報表管線、每月的備份管線,應該拆成三個 DAG,彼此透過 Asset 或 ExternalTaskSensor 連接。一個 DAG 跑太久(超過幾小時)會拖累整個排程器,拆分是必要的。
Backfill 與 Catchup
Backfill 是「補跑」的概念。當你的 DAG 從 11 月 1 日開始部署,但 11 月 5 日才上線,你想把 1 日到 4 日的排程都補跑,Airflow 提供兩種方式:
- catchup=True:DAG 一啟動就自動補跑從 start_date 到現在的所有排程。預設 False,建議新手先保持 False。
- airflow dags backfill:手動觸發補跑,可以指定日期區間。
airflow dags backfill --start-date 2025-11-01 --end-date 2025-11-04 daily_orders。
Backfill 對增量模型(incremental)特別有用。當你改了 dbt 模型,想要用新邏輯重算歷史三個月的資料,Backfill 可以讓 DAG 自動依序補跑每一天的排程。
連線(Connection)與變數(Variable)
Airflow 把外部系統的連線資訊(資料庫、API 金鑰)集中管理:
- Connection:資料庫、S3、HTTP 端點的連線字串。在 Airflow UI 或
airflow connections addCLI 建立。Task 透過conn_id引用。 - Variable:任意 key-value 設定,例如
batch_size=10000。在 UI 或airflow variables setCLI 建立。Task 透過Variable.get("batch_size")讀取。
把連線資訊集中管理的好處是「換環境不必改程式」。同一支 DAG 在 dev 環境用 dev 資料庫、在 prod 環境用 prod 資料庫,靠 Connection ID 自動切換。
常見錯誤與踩雷
錯誤一:DAG 檔案裡有 top-level code 過於昂貴。Airflow 排程器會定期掃描所有 DAG 檔案(預設每 5 分鐘),如果 DAG 檔案開頭就呼叫 API、抓大檔案,整個排程器會被拖慢。原則:DAG 檔案裡只能定義 Task,不能執行副作用。
錯誤二:Task 之間的 XCom 傳遞大資料。XCom 是設計來傳小設定(檔名、ID、字串)的,不是設計來傳整個 DataFrame。如果你發現自己在 XCom 裡塞了 100MB 的資料,那個架構已經錯了。應該改成把資料落地到 S3 或 DuckDB,只在 XCom 傳檔名。
錯誤三:把所有排程都擠在整點。如果 5 條 DAG 都設定 schedule="0 2 * * *"(每天凌晨兩點),Airflow 會同時跑 5 條管線,worker 可能不夠。實務上會把排程錯開幾分鐘(02:00、02:05、02:10)讓資源平均分配。
錯誤四:用 cron 表達式但忘記時區。Airflow 預設使用 UTC。台北時間凌晨兩點 = UTC 18:00。寫成 schedule="0 2 * * *" 會被當成 UTC 02:00 = 台北 10:00。如果要鎖台北時間,用 schedule="0 18 * * *",或是在 start_date 用 timezone("Asia/Taipei")。
錯誤五:忘了設定 retries。預設 retries=0,意思是 Task 失敗一次就停止。對接外部 API、抓網路資源的 Task,建議 retries=3 與 retry_delay=timedelta(minutes=5)。
效能與實務提醒
Airflow 3.x 的執行效率比 2.x 提升不少,主要來自 Task SDK 的解耦設計。但有幾個調校重點:
- parallelism 設定:預設 32 個 Task 同時執行。如果 worker 比較少可以調低。
- scheduler 與 worker 分離:Airflow 3.x 預設把 scheduler 與 executor 分離,scheduler 只負責排程,executor 負責執行 Task。
- dag_file_processor_timeout:DAG 檔案處理超時時間,預設 50 秒。當 DAG 很多時可以調長。
實務上,當管線數量成長到 30 條以上,你會開始感覺 Airflow 變重——這時候就是考慮 Day 29 的輕量替代方案的時機。但在那之前,Airflow 的標準化與生態系優勢依然明顯。
小結
今天我們用「看圖說故事」的方式走過 Airflow 3.x 的核心觀念:DAG 是「任務的有向無環圖」,Task 是最小執行單位,Operator 與 Sensor 是 Task 的型別,TaskFlow API 是 3.x 推薦的 Python 寫法,Asset 是 3.x 引入的「資料驅動排程」概念。我們也示範了 DAG 的設計原則(冪等、顆粒度、錯誤處理)與常見陷阱(top-level code、XCom 大資料、時區)。
這些觀念明天會被落實成實際的 Docker 部署。你會看到一份 docker-compose.yaml,把 Airflow 3.x 連同 Postgres 與 Redis 一起拉起來,並把昨天的 dbt 管線掛上去。
結語
Airflow 是資料工程的「編排作業系統」——一旦裝起來、用熟,所有資料管線都會走同一套流程。但它的學習曲線與系統需求都不低,因此理解觀念比急著部署更重要。今天這篇先把概念釐清,明天 Day 28 我們再實際拉一支 Docker 版本,把昨天 Day 25 的 dbt 專案掛進去。
如果你是那種「Airflow 對我來說太重了」的讀者,明天看完 Docker 版本後,可以跳到 Day 29 的輕量替代方案——用 APScheduler、Prefect、Dagster 或純 Python 排程,照樣能達成 80% 的效果,但學習成本低很多。
延伸資源
- Apache Airflow 官方說明「Tutorial」與「Core Concepts」段落,本篇所有術語的權威來源。
- 《Data Pipelines Pocket Reference》by James Densmore(O'Reilly 2021),簡明扼要地介紹 Airflow 與替代工具的取捨。
- Airflow 3.x 官方公告與 migration guide,理解從 2.x 升級 3.x 的關鍵改動。
留言
張貼留言