Dynamisches Task-Mapping

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Das Problem mit hart codierten Tasks

@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")

Beispiel für dynamisches Task-Mapping

Data-Pipelines mit Airflow aufbauen

.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) # maximale Parallelität steuern def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

Zur Laufzeit skalieren

Data-Pipelines mit Airflow aufbauen

Gemappte Tasks in der UI

 

  • Gemappte Instanzen erscheinen als indizierte Instanzen unter einem einzelnen Task
  • Klicke in jede Instanz für eigene Logs und Status
  • Der Dag bleibt übersichtlich, selbst mit Hunderten Instanzen

Gemappte Tasks in der UI

Data-Pipelines mit Airflow aufbauen

.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 und expand

Data-Pipelines mit Airflow aufbauen

Das Kreuzprodukt-Problem

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

# Kreuzprodukt: 3 x 3 = 9 Instanzen (nicht gewünscht!) process.expand(path=files, dest=destinations)

Listen zippen

Data-Pipelines mit Airflow aufbauen

.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"]

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

 

  • zip() paart Elemente eins-zu-eins nach Position
  • expand_kwargs entpackt jedes Paar in Keyword-Argumente
Data-Pipelines mit Airflow aufbauen

Wann was verwenden

 

Pattern Use case
expand() Über eine einzelne Liste mappen
partial() Gemeinsame Parameter für alle Instanzen festlegen
zip() Mehrere Listen positionsweise paaren

 

  • Hilft beim Aufbau skalierbarer Dags
  • Wenn sich Daten ändern, passt sich die Taskanzahl zur Laufzeit automatisch an
Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...