Organizzare DAG complessi con i Task Group

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Il problema della complessità

Dag complesso

  • I DAG grandi diventano difficili da navigare nella vista Graph
  • I nomi dei task si confondono
  • I nuovi membri faticano a trovare ciò di cui sono responsabili
Creare data pipeline con 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 imposta un identificatore personalizzato, per default è il nome della funzione
  • default_args si applica a tutti i task del gruppo per evitare codice di configurazione ripetuto
  • I task group appaiono come blocchi espandibili nella UI
Creare data pipeline con Airflow

Come appare nella UI di Airflow

Senza task group

Grafico semplice con tre task in linea: extract_orders, transform_orders, load, tutti allo stesso livello

  • Tutti i task allo stesso livello

Con task group

Grafico che mostra un blocco compresso chiamato process_orders con extract e transform, seguito da un task load fuori dal gruppo

  • Blocchi comprimibili nella UI
Creare data pipeline con Airflow

Annidamento e nomi visualizzati personalizzati

@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(),
    }
  • I task group possono essere annidati per pipeline più complesse
  • group_display_name imposta un etichetta leggibile nella UI (supporta anche emoji)
Creare data pipeline con Airflow

Il factory pattern

@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")
  • I task group sono funzioni Python decorate
  • Puoi chiamarli più volte con parametri diversi
  • Questo factory pattern consente di riusare la logica della pipeline
Creare data pipeline con Airflow

Linee guida per il raggruppamento

Raggruppare per dominio o ambito

  • Raggruppa per dominio o ambito, non per tipo di operatore
  • Usa group_id per nomi programmatici chiari, group_display_name per le etichette UI
  • Usa default_args per condividere config come i retries tra tutti i task del gruppo
  • Applica la Legge di Miller: più di 7 elementi top-level può indicare che serve un task group
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...