Organizowanie złożonych DAG-ów za pomocą grup zadań

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problem ze złożonością

Complex Dag

  • Duże DAG-i są trudne w nawigacji w widoku grafu
  • Nazwy zadań zlewają się ze sobą
  • Nowi członkowie zespołu mają problem ze znalezieniem swoich zadań
Budowanie potoków danych z 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 ustawia własny identyfikator — domyślnie przyjmuje nazwę funkcji
  • default_args stosuje się do wszystkich zadań w grupie, eliminując powtarzający się kod konfiguracyjny
  • Grupy zadań wyświetlają się w interfejsie jako bloki, które można zwijać i rozwijać
Budowanie potoków danych z Airflow

Jak to wygląda w interfejsie Airflow

Bez grup zadań

Prosty graf z trzema zadaniami w linii: extract_orders, transform_orders, load — wszystkie na tym samym poziomie

  • Wszystkie zadania na tym samym poziomie

Z grupami zadań

Graf ze zwiniętym blokiem process_orders zawierającym zadania extract i transform, po którym następuje zadanie load poza grupą

  • Zwijane bloki w interfejsie
Budowanie potoków danych z Airflow

Zagnieżdżanie i własne nazwy wyświetlane

@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(),
    }
  • Grupy zadań można zagnieżdżać w bardziej złożonych potokach
  • group_display_name ustawia czytelną etykietę w interfejsie (obsługuje też emoji)
Budowanie potoków danych z Airflow

Wzorzec fabryki

@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")
  • Grupy zadań to dekorowane funkcje Pythona
  • Możesz je wywoływać wielokrotnie z różnymi parametrami
  • Ten wzorzec fabryki pozwala wielokrotnie używać tej samej logiki potoku
Budowanie potoków danych z Airflow

Zasady grupowania

Group by domain or concern

  • Grupuj według domeny lub obszaru odpowiedzialności, a nie według typu operatora
  • Używaj group_id do nazw programistycznych, a group_display_name do etykiet w interfejsie
  • Używaj default_args, by współdzielić konfigurację, np. liczbę ponownych prób, w całej grupie
  • Stosuj prawo Milleraponad 7 elementów na najwyższym poziomie to sygnał, że przyda się grupa zadań
Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...