पाठ 15 / 25
Dynamic Task Mapping
`expand()` से run के समय परिवर्तनीय संख्या में समानांतर tasks बनाएँ।
एक परिभाषा, कई tasks
कभी आपको तब ही पता चलता है कि कितने समानांतर काम चाहिए जब DAG चलता है, जैसे "आज आई हर फ़ाइल प्रोसेस करो"। Dynamic task mapping task को सूची पर expand होने देता है ताकि Airflow run के समय हर item के लिए एक mapped task instance बनाए। बदलते argument के लिए .expand(arg=list_or_xcom) और निश्चित arguments के लिए .partial(...) उपयोग करें। Downstream task तब सारे नतीजे सूची के रूप में पा सकता है। Mapped instances की संख्या उचित रखें (सैकड़ों, लाखों नहीं)।
सूची पर expand
list_files task नाम लौटाता है; process.expand हर नाम के लिए एक task बनाता है; summarise सारे return values पाता है। मैंने इसे Airflow 3.0.6 पर airflow dags test process_files 2026-10-01 से चलाया; तीन mapped process instances चले और summarise ने नीचे की पंक्ति छापी।
from datetime import datetime
from airflow.sdk import dag, task
@dag(schedule="@daily", start_date=datetime(2026, 9, 1), catchup=False)
def process_files():
@task
def list_files() -> list[str]:
return ["a.csv", "b.csv", "c.csv"]
@task
def process(name: str) -> int:
return len(name) # stand-in for real work
@task
def summarise(counts: list[int]) -> None:
print(f"processed {len(counts)} files")
summarise(process.expand(name=list_files()))
process_files()
Output:
processed 3 files
त्वरित जाँच: Dynamic task mapping कब उपयोगी है?
- जब समानांतर tasks की संख्या सिर्फ़ run के समय पता हो
- जब आप DAG फ़ाइल चलते समय बदलना चाहें
- जब कोई task न हो
- कभी नहीं
Answer
जब समानांतर tasks की संख्या सिर्फ़ run के समय पता हो — Mapping run के दौरान बने डेटा से task instances बनाता है।