Emparejar entradas con zip
Tu equipo tiene un pipeline que sube archivos a ubicaciones específicas de almacenamiento en la nube. Una persona del equipo usa zip() con expand_kwargs() para emparejar cada archivo con su destino en la nube:
@task
def get_files():
return ["users.csv", "orders.csv", "products.csv"]
@task
def get_destinations():
return ["s3://bronze/users/", "s3://bronze/orders/", "s3://bronze/products/"]
@task
def upload(file_to_upload):
# Simulate uploading a file to cloud storage
print(f"Uploading {file_to_upload[0]} to {file_to_upload[1]}")
files = get_files()
destinations = get_destinations()
upload.expand(file_to_upload=files.zip(destinations))
¿Cuántas instancias de upload creará Airflow?
Este ejercicio forma parte del curso
Creación de canalizaciones de datos con Airflow
ejercicio interactivo práctico
Convierte la teoría en práctica con uno de nuestros ejercicios interactivos
Empezar ejercicio