पाठ 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 बनाता है।