Complexe Dags organiseren met Task Groups

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Het complexiteitsprobleem

Complexe Dag

  • Grote Dags worden lastig te navigeren in de Graph-weergave
  • Task-namen vloeien in elkaar over
  • Nieuwe teamleden vinden lastig wat van hen is
Data-pijplijnen bouwen met 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 stelt een aangepaste id in, standaard de functienaam
  • default_args geldt voor alle taken in de groep, zodat je configuratie niet hoeft te herhalen
  • Task groups verschijnen als uitklapbare blokken in de UI
Data-pijplijnen bouwen met Airflow

Zo ziet het eruit in de Airflow UI

Zonder task groups

Eenvoudige graaf met drie taken op een rij: extract_orders, transform_orders, load, allemaal op hetzelfde niveau

  • Alle taken op hetzelfde niveau

Met task groups

Graaf met een ingeklapt blok genaamd process_orders met extract en transform, gevolgd door een load-taak buiten de groep

  • Inklapbare blokken in de UI
Data-pijplijnen bouwen met Airflow

Nesten en aangepaste weergavenamen

@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 groups kun je nesten voor complexere pipelines
  • group_display_name stelt een leesbaar label in de UI in (ondersteunt zelfs emoji's)
Data-pijplijnen bouwen met Airflow

Het 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())

# Hergebruik hetzelfde patroon voor verschillende bronnen orders = process_source("orders", "/data/orders.csv") returns = process_source("returns", "/data/returns.csv") events = process_source("events", "/data/events.csv")
  • Task groups zijn gedecoreerde Python-functies
  • Je kunt ze meerdere keren aanroepen met andere parameters
  • Dit factory pattern laat je pijplogica hergebruiken
Data-pijplijnen bouwen met Airflow

Richtlijnen voor groeperen

Groepeer op domein of concern

  • Groepeer op domein of concern, niet op operatortype
  • Gebruik group_id voor programmeerbare namen, group_display_name voor UI-labels
  • Gebruik default_args om config zoals retries te delen binnen een groep
  • Pas Wet van Miller toe, meer dan 7 items op topniveau kan betekenen dat je een task group nodig hebt
Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...