Airflow के साथ Data Pipelines बनाना
Volker Janz
Senior Developer Advocate at Astronomer

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 के

task groups के साथ

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

group_id और UI labels के लिए group_display_name उपयोग करेंdefault_args उपयोग करेंAirflow के साथ Data Pipelines बनाना