EmpezarEmpieza gratis

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

Ver curso

ejercicio interactivo práctico

Convierte la teoría en práctica con uno de nuestros ejercicios interactivos

Empezar ejercicio