Mapeo dinámico de tareas

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

El problema de las tareas con valores fijos

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

Ejemplo de Dynamic Task Mapping

Creación de canalizaciones de datos 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) # controlar la paralelización máxima def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

Escala en tiempo de ejecución

Creación de canalizaciones de datos con Airflow

Tareas mapeadas en la interfaz

 

  • Las instancias mapeadas aparecen como instancias indexadas bajo una sola tarea
  • Entra en cada una para ver logs y estado individuales
  • El DAG sigue siendo legible incluso con cientos de instancias

Tareas mapeadas en la interfaz

Creación de canalizaciones de datos 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 y expand

Creación de canalizaciones de datos con Airflow

El problema del producto 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"]

# Producto cartesiano: 3 x 3 = 9 instancias (¡no es lo que queremos!) process.expand(path=files, dest=destinations)

Unir listas con zip

Creación de canalizaciones de datos 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"]

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

 

  • zip() empareja elementos uno a uno por posición
  • expand_kwargs descomprime cada par en argumentos con nombre
Creación de canalizaciones de datos con Airflow

Cuándo usar cada uno

 

Pattern Use case
expand() Mapear sobre una lista única
partial() Fijar parámetros compartidos entre instancias
zip() Emparejar varias listas por posición

 

  • Ayudan a construir DAGs escalables
  • Cuando cambian los datos, el número de tareas se ajusta automáticamente en tiempo de ejecución
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...