Organisera komplexa DAG:ar med uppgiftsgrupper

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Komplexitetsproblemet

Komplex DAG

  • Stora DAG:ar blir svåra att navigera i grafvyn
  • Uppgiftsnamnen flyter ihop
  • Nya teammedlemmar har svårt att hitta sina uppgifter
Bygg datapipelines med 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 anger ett eget identifieringsnamn, och används som standard funktionens namn
  • default_args tillämpas på alla uppgifter i gruppen och minskar upprepad konfigurationskod
  • Uppgiftsgrupper visas som expanderbara block i gränssnittet
Bygg datapipelines med Airflow

Så ser det ut i Airflow-gränssnittet

Utan uppgiftsgrupper

Enkel graf med tre uppgifter på rad: extract_orders, transform_orders och load, alla på samma nivå

  • Alla uppgifter på samma nivå

Med uppgiftsgrupper

Graf med ett ihopfällt block kallat process_orders som innehåller extract och transform, följt av en load-uppgift utanför gruppen

  • Hopfällbara block i gränssnittet
Bygg datapipelines med Airflow

Nästlade grupper och anpassade visningsnamn

@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(),
    }
  • Uppgiftsgrupper kan nästlas för mer komplexa pipelines
  • group_display_name anger en läsbar etikett i gränssnittet (även emojis stöds)
Bygg datapipelines med Airflow

Fabriksmönstret

@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")
  • Uppgiftsgrupper är dekorerade Python-funktioner
  • Du kan anropa dem flera gånger med olika parametrar
  • Det här fabriksmönstret gör det möjligt att återanvända pipelinelogik
Bygg datapipelines med Airflow

Riktlinjer för gruppering

Gruppera efter domän eller ansvarsområde

  • Gruppera efter domän eller ansvarsområde, inte efter operatortyp
  • Använd group_id för tydliga programmatiska namn och group_display_name för gränssnittsetiketter
  • Använd default_args för att dela konfiguration som återförsök över alla uppgifter i en grupp
  • Tillämpa Millers lagfler än 7 toppnivåobjekt kan vara ett tecken på att du behöver en uppgiftsgrupp
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...