Task Groups के साथ जटिल Dags को व्यवस्थित करना

Airflow के साथ Data Pipelines बनाना

Volker Janz

Senior Developer Advocate at Astronomer

जटिलता की समस्या

जटिल Dag

  • बड़े Dags Graph view में नेविगेट करना मुश्किल हो जाते हैं
  • टास्क नाम आपस में घुल जाते हैं
  • नए टीम सदस्य अपनी ओनरशिप पहचानने में जूझते हैं
Airflow के साथ Data Pipelines बनाना

@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 एक कस्टम पहचानकर्ता सेट करता है, डिफॉल्ट फ़ंक्शन नाम होता है
  • default_args ग्रुप के सभी tasks पर लागू होते हैं ताकि दोहराया कॉन्फ़िग कोड न लिखना पड़े
  • Task groups UI में expandable ब्लॉक्स की तरह दिखते हैं
Airflow के साथ Data Pipelines बनाना

Airflow UI में यह कैसा दिखता है

बिना task groups के

एक सरल ग्राफ़ जिसमें तीन टास्क लाइन में: extract_orders, transform_orders, load, सभी एक ही स्तर पर

  • सभी टास्क एक ही स्तर पर

task groups के साथ

ग्राफ़ जिसमें process_orders नाम का एक collapsed ब्लॉक है जिसमें extract और transform हैं, और उसके बाद बाहर एक load टास्क है

  • UI में collapsible ब्लॉक्स
Airflow के साथ Data Pipelines बनाना

Nesting और custom display names

@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 को nested कर सकते हैं
  • group_display_name UI में मानव-पठनीय लेबल सेट करता है (emojis भी सपोर्टेड)
Airflow के साथ Data Pipelines बनाना

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 groups decorated Python functions होते हैं
  • आप उन्हें अलग-अलग parameters के साथ कई बार कॉल कर सकते हैं
  • यह factory pattern पाइपलाइन लॉजिक को रीयूज़ करने देता है
Airflow के साथ Data Pipelines बनाना

Grouping के लिए दिशानिर्देश

डोमेन या concern के आधार पर समूह बनाना

  • डोमेन या concern के आधार पर समूह बनाएँ, operator टाइप के आधार पर नहीं
  • प्रोग्रामेटिक स्पष्ट नामों के लिए group_id और UI labels के लिए group_display_name उपयोग करें
  • ग्रुप में सभी tasks के लिए retries जैसी कॉन्फ़िग साझा करने हेतु default_args उपयोग करें
  • Miller's Law अपनाएँ, 7 से अधिक top-level items हों तो task group की ज़रूरत का संकेत है
Airflow के साथ Data Pipelines बनाना

अभ्यास करते हैं!

Airflow के साथ Data Pipelines बनाना

Preparing Video For Download...