Mappatura dinamica dei task

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Il problema dei task hardcoded

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

Esempio di Dynamic Task Mapping

Creare data pipeline con 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) # controlla la parallelizzazione massima def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

Scala a runtime

Creare data pipeline con Airflow

Task mappati nella UI

 

  • Le istanze mappate compaiono come istanze indicizzate sotto un unico task
  • Clicca su ciascuna per log e stato individuali
  • Il DAG resta leggibile anche con centinaia di istanze

Task mappati nella UI

Creare data pipeline con 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 ed expand

Creare data pipeline con Airflow

Il problema del prodotto cartesiano

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

# Prodotto cartesiano: 3 x 3 = 9 istanze (non è ciò che vogliamo!) process.expand(path=files, dest=destinations)

Unisci liste con zip

Creare data pipeline con 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"]

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

 

  • zip() abbina gli elementi uno a uno in base alla posizione
  • expand_kwargs espande ogni coppia in argomenti keyword
Creare data pipeline con Airflow

Quando usare cosa

 

Pattern Caso d'uso
expand() Mappare su una singola lista
partial() Fissare parametri condivisi tra le istanze
zip() Accoppiare più liste per posizione

 

  • Aiuta a creare DAG scalabili
  • Quando i dati cambiano, il numero di task si adatta automaticamente a runtime
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...