पाठ 6 / 25

TaskFlow API

@dag और @task से DAGs लिखें ताकि Python functions tasks बनें और return values XComs बनें।

सादे functions को tasks

TaskFlow API tasks को @task decorator वाले सामान्य Python functions के रूप में लिखने देता है और एक को दूसरे का output देकर उन्हें जोड़ता है। Airflow return value को XCom के रूप में रखता है और निर्भरता अनुमान से बना लेता है, इसलिए load(transform(extract())) extract → transform → load श्रृंखला बनाता है। यह संक्षिप्त है और Python-आधारित tasks के लिए सुझाई शैली है।

TaskFlow से छोटा ETL

उदाहरण को आत्मनिर्भर रखने के लिए डेटा hard-coded है। असली pipeline में extract API या database बुलाता, और tasks के बीच सिर्फ़ छोटे मान जाते। मैंने इसे Airflow 3.0.6 पर airflow dags test orders_etl 2026-10-01 से चलाया; load task ने कुल छापा, 120 + 80 + 200 = 400।

from datetime import datetime

from airflow.sdk import dag, task


@dag(
    dag_id="orders_etl",
    start_date=datetime(2026, 9, 1),
    schedule="@daily",
    catchup=False,
    default_args={"retries": 2},
    tags=["demo"],
)
def orders_etl():
    @task
    def extract() -> list[dict]:
        return [{"id": 1, "amount": 120}, {"id": 2, "amount": 80}, {"id": 3, "amount": 200}]

    @task
    def transform(rows: list[dict]) -> int:
        return sum(r["amount"] for r in rows)

    @task
    def load(total: int) -> None:
        print(f"Total revenue loaded: {total}")

    load(transform(extract()))


orders_etl()

Output:

Total revenue loaded: 400

Task functions में type hints दें

Return type hints दिखाते हैं कि tasks के बीच क्या बहता है और ग़लतियाँ जल्दी पकड़ने में मदद करते हैं। मान छोटे और serialise होने योग्य रखें (संख्याएँ, strings, छोटी lists और dicts)।

त्वरित जाँच: TaskFlow में एक @task से अगले तक डेटा कैसे जाता है?

  • डेटा नहीं जा सकता
  • Global variable से
  • Email से
  • Return value के रूप में, XCom में रखा हुआ
Answer

Return value के रूप में, XCom में रखा हुआ — Return values XComs बनते हैं और Airflow निर्भरता अपने आप जोड़ता है।