Komplexe DAGs mit Task-Gruppen organisieren

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Das Komplexitätsproblem

Komplexer Dag

  • Große DAGs sind in der Graph-Ansicht schwer zu navigieren
  • Aufgabennamen verschwimmen
  • Neue Teammitglieder finden ihren Bereich kaum
Data-Pipelines mit Airflow aufbauen

@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 setzt eine individuelle Kennung, standardmäßig der Funktionsname
  • default_args gilt für alle Tasks der Gruppe und vermeidet doppelten Konfig-Code
  • Task-Gruppen erscheinen als aufklappbare Blöcke in der UI
Data-Pipelines mit Airflow aufbauen

So sieht es in der Airflow-UI aus

Ohne Task-Gruppen

Einfacher Graph mit drei Tasks in einer Linie: extract_orders, transform_orders, load, alle auf derselben Ebene

  • Alle Tasks auf derselben Ebene

Mit Task-Gruppen

Graph mit einem eingeklappten Block namens process_orders mit extract und transform, gefolgt von einem load-Task außerhalb der Gruppe

  • Auf- und zuklappbare Blöcke in der UI
Data-Pipelines mit Airflow aufbauen

Verschachteln und Anzeigenamen anpassen

@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-Gruppen lassen sich für komplexere Pipelines verschachteln
  • group_display_name setzt eine lesbare Bezeichnung in der UI (unterstützt sogar Emojis)
Data-Pipelines mit Airflow aufbauen

Das 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")
  • Task-Gruppen sind dekorierte Python-Funktionen
  • Du kannst sie mit unterschiedlichen Parametern mehrfach aufrufen
  • Dieses Factory-Pattern ermöglicht Wiederverwendung von Pipeline-Logik
Data-Pipelines mit Airflow aufbauen

Richtlinien fürs Gruppieren

Nach Domäne oder Concern gruppieren

  • Nach Domäne oder Concern gruppieren, nicht nach Operator-Typ
  • group_id für klare programmgesteuerte Namen, group_display_name für UI-Labels
  • default_args für geteilte Config wie retries über alle Tasks der Gruppe
  • Miller's Law beachten: mehr als 7 Top-Level-Elemente deuten auf Bedarf für eine Task-Gruppe hin
Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...