पाठ 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 निर्भरता अपने आप जोड़ता है।