Dynamische task mapping

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Het probleem met hardgecodeerde taken

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

Voorbeeld van Dynamic Task Mapping

Data-pijplijnen bouwen met Airflow

.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) # max. parallelisme beperken def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

Schaal op runtime

Data-pijplijnen bouwen met Airflow

Gemapte taken in de UI

 

  • Gemapte instanties verschijnen als geïndexeerde instanties onder één taak
  • Klik per instantie voor aparte logs en status
  • Dag blijft leesbaar, ook met honderden instanties

Gemapte taken in de UI

Data-pijplijnen bouwen met Airflow

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

Data-pijplijnen bouwen met Airflow

Het kruisproductprobleem

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

# Cartesisch product: 3 x 3 = 9 instanties (niet wat we willen!) process.expand(path=files, dest=destinations)

Lijsten zippen

Data-pijplijnen bouwen met Airflow

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

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

 

  • zip() koppelt items één-op-één op positie
  • expand_kwargs pakt elk paar uit tot keyword-argumenten
Data-pijplijnen bouwen met Airflow

Wanneer gebruik je wat

 

Patroon Gebruik
expand() Mappen over één lijst
partial() Gedeelde parameters vastzetten over instanties
zip() Meerdere lijsten paren op positie

 

  • Helpen bij het bouwen van schaalbare Dags
  • Als de data verandert, past het aantal taken zich automatisch aan tijdens runtime aan
Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...