Organizarea DAG-urilor complexe cu Task Groups

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problema complexității

DAG complex

  • DAG-urile mari devin greu de navigat în vizualizarea Graph
  • Numele sarcinilor se confundă
  • Membrii noi ai echipei nu știu ce le aparține
Construirea pipeline-urilor de date cu 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 setează un identificator personalizat, implicit numele funcției
  • default_args se aplică tuturor sarcinilor din grup, eliminând configurarea repetitivă
  • Task group-urile apar ca blocuri expandabile în UI
Construirea pipeline-urilor de date cu Airflow

Cum arată în Airflow UI

Fără task group-uri

Graf simplu cu trei sarcini în linie: extract_orders, transform_orders, load, toate la același nivel

  • Toate sarcinile la același nivel

Cu task group-uri

Graf cu un bloc restrâns numit process_orders, conținând extract și transform, urmat de o sarcină load în afara grupului

  • Blocuri reductibile în UI
Construirea pipeline-urilor de date cu Airflow

Imbricare și nume de afișare personalizate

@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-urile pot fi imbricate pentru pipeline-uri mai complexe
  • group_display_name setează o etichetă lizibilă în UI (acceptă și emoji)
Construirea pipeline-urilor de date cu Airflow

Pattern-ul factory

@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-urile sunt funcții Python decorate
  • Pot fi apelate de mai multe ori cu parametri diferiți
  • Acest pattern factory permite reutilizarea logicii din pipeline
Construirea pipeline-urilor de date cu Airflow

Recomandări pentru grupare

Grupare după domeniu sau responsabilitate

  • Grupează după domeniu sau responsabilitate, nu după tipul operatorului
  • Folosește group_id pentru nume programatice clare, group_display_name pentru etichete în UI
  • Folosește default_args pentru a partaja configurații precum reîncercările între toate sarcinile din grup
  • Aplică Legea lui Miller: mai mult de 7 elemente de nivel superior poate indica necesitatea unui task group
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...