用 Task Group 組織複雜的 Dag

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

複雜度問題

複雜的 Dag

  • 大型 Dag 在 Graph 檢視中難以瀏覽
  • 任務名稱擠在一起
  • 新成員難以找到自己的負責範圍
使用 Airflow 建置資料管線

@task_group

from airflow.sdk import dag, task, task_group

@task_group(
    group_id="ingest_orders",
    default_args={"retries": 3},
)

def process_orders(): @task def extract_orders(): return [{"id": 1, "amount": 99.99}] @task def transform_orders(orders): return [{"id": o["id"], "total": o["amount"] * 1.08} for o in orders] return transform_orders(extract_orders())
  • group_id 設定自訂識別碼,預設為函式名稱
  • default_args 套用到群組內所有任務,避免重複的設定程式碼
  • Task group 會在 UI 中顯示為可展開的區塊
使用 Airflow 建置資料管線

在 Airflow UI 中的樣子

未使用 task group

簡單圖,三個同層任務串聯:extract_orders、transform_orders、load

  • 所有任務在同一層級

使用 task group

圖示:一個名為 process_orders 的摺疊區塊,內含 extract 與 transform,後接群組外的 load 任務

  • UI 中可摺疊的區塊
使用 Airflow 建置資料管線

巢狀與自訂顯示名稱

@task_group(group_display_name="Process All")
def process_all():

    @task_group(group_display_name="Ingest Orders")
    def orders():
        return transform(extract())

    @task_group(group_display_name="Process Returns")
    def returns():
        return transform(extract())

    return {
        "orders": orders(),
        "returns": returns(),
    }
  • Task group 可巢狀,用於更複雜的 pipeline
  • group_display_name 設定 UI 中易讀的標籤(支援表情符號)
使用 Airflow 建置資料管線

工廠模式

@task_group
def process_source(source_name, source_path):

    @task
    def extract():
        return read_data(source_path)

    @task
    def transform(data):
        return clean_data(data)

    return transform(extract())

# Reuse the same pattern for different sources orders = process_source("orders", "/data/orders.csv") returns = process_source("returns", "/data/returns.csv") events = process_source("events", "/data/events.csv")
  • Task group 是加上裝飾器的 Python 函式
  • 你可以用不同參數重複呼叫
  • 這種工廠模式重用 pipeline 邏輯
使用 Airflow 建置資料管線

分組準則

依網域或關注點分組

  • 網域或關注點分組,而非依 operator 類型
  • group_id 用於清楚的程式化名稱group_display_name 用於UI 標籤
  • 使用 default_args 共用如重試等設定到整個群組
  • 套用米勒定律超過 7 個頂層項目可能表示需要 task group
使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...