跳到主要內容

DE Day 27 Airflow 概念與 DAG 設計

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 add CLI 建立。Task 透過 conn_id 引用。
  • Variable:任意 key-value 設定,例如 batch_size=10000。在 UI 或 airflow variables set CLI 建立。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 的解耦設計。但有幾個調校重點:

  1. parallelism 設定:預設 32 個 Task 同時執行。如果 worker 比較少可以調低。
  2. scheduler 與 worker 分離:Airflow 3.x 預設把 scheduler 與 executor 分離,scheduler 只負責排程,executor 負責執行 Task。
  3. 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 的關鍵改動。

留言

這個網誌中的熱門文章

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 建構深度學習模型。 開發者與研究人員 :想更深入了...