Organizace složitých DAGů pomocí skupin úkolů

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problém se složitostí

Složitý DAG

  • Velké DAGy se v zobrazení grafu špatně orientují
  • Názvy úkolů splývají
  • Noví členové týmu těžko hledají, co mají na starosti
Tvorba datových pipeline s 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 nastavuje vlastní identifikátor, výchozí hodnotou je název funkce
  • default_args se vztahuje na všechny úkoly ve skupině a zamezuje opakování konfigurace
  • Skupiny úkolů se v UI zobrazují jako rozbalovací bloky
Tvorba datových pipeline s Airflow

Jak to vypadá v Airflow UI

Bez skupin úkolů

Jednoduchý graf se třemi úkoly v řadě: extract_orders, transform_orders, load – všechny na stejné úrovni

  • Všechny úkoly na stejné úrovni

Se skupinami úkolů

Graf zobrazující sbalený blok process_orders s úkoly extract a transform, za nímž následuje úkol load mimo skupinu

  • Sbalitelné bloky v UI
Tvorba datových pipeline s Airflow

Vnořování a vlastní zobrazované názvy

@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(),
    }
  • Skupiny úkolů lze vnořovat pro složitější pipeline
  • group_display_name nastavuje čitelný popisek v UI (podporuje i emoji)
Tvorba datových pipeline s Airflow

Tovární vzor

@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")
  • Skupiny úkolů jsou dekorované Python funkce
  • Lze je volat vícekrát s různými parametry
  • Tento tovární vzor umožňuje znovupoužití logiky pipeline
Tvorba datových pipeline s Airflow

Doporučení pro seskupování

Seskupování podle domény nebo odpovědnosti

  • Seskupuj podle domény nebo odpovědnosti, ne podle typu operátoru
  • Používej group_id pro jasné programové názvy a group_display_name pro popisky v UI
  • Pomocí default_args sdílej konfiguraci, například počet opakování, pro všechny úkoly ve skupině
  • Uplatňuj Millerův zákonvíce než 7 položek na nejvyšší úrovni může být signál, že je čas na skupinu úkolů
Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...