डायनेमिक टास्क मैपिंग

Airflow के साथ Data Pipelines बनाना

Volker Janz

Senior Developer Advocate at Astronomer

हार्डकोडेड टास्क की समस्या

@task
def process_file_1():
    transform("/data/file_1.csv")
@task
def process_file_2():
    transform("/data/file_2.csv")
@task
def process_file_3():
    transform("/data/file_3.csv")

डायनेमिक टास्क मैपिंग उदाहरण

Airflow के साथ Data Pipelines बनाना

.expand()

@task
def fetch_files():
    return ["/data/file_1.csv", "/data/file_2.csv", "/data/file_3.csv"]

@task(max_active_tis_per_dagrun=2) # control max parallelism def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

रनटाइम पर स्केल करें

Airflow के साथ Data Pipelines बनाना

UI में मैप्ड टास्क

 

  • मैप की गई इंस्टैंसेज़ एक ही टास्क के अंतर्गत इंडेक्स्ड इंस्टैंसेज़ के रूप में दिखती हैं
  • हर एक में क्लिक करके अलग लॉग और स्टेटस देखें
  • सैकड़ों इंस्टैंसेज़ होने पर भी Dag पढ़ने में आसान रहता है

UI में मैप्ड टास्क

Airflow के साथ Data Pipelines बनाना

.partial()

@task
def process_file(path: str, output_dir: str):
    transform(path, output_dir)

files = get_files() process_file.partial(output_dir="/out").expand(path=files)

Partial और expand

Airflow के साथ Data Pipelines बनाना

क्रॉस-प्रोडक्ट समस्या

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# Cross product: 3 x 3 = 9 instances (not what we want!) process.expand(path=files, dest=destinations)

लिस्ट्स को zip करें

Airflow के साथ Data Pipelines बनाना

.zip()

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# Paired: 3 instances (a->a, b->b, c->c) process.expand_kwargs(files.zip(destinations))

 

  • zip() पोज़िशन से आइटम्स को एक-से-एक मिलाता है
  • expand_kwargs हर पेयर को keyword arguments में अनपैक करता है
Airflow के साथ Data Pipelines बनाना

क्या कब उपयोग करें

 

पैटर्न उपयोग
expand() एक single list पर मैप करें
partial() इंस्टैंसेज़ में shared parameters फिक्स करें
zip() कई लिस्ट्स को पोज़िशन से पेयर करें

 

  • स्केलेबल Dags बनाने में मदद करें
  • डेटा बदलने पर टास्क काउंट रनटाइम पर स्वतः एडजस्ट होता है
Airflow के साथ Data Pipelines बनाना

अभ्यास करते हैं!

Airflow के साथ Data Pipelines बनाना

Preparing Video For Download...