Inizia subitoInizia gratis

Abbinare gli input con zip

Il tuo team ha una pipeline che carica file in specifiche destinazioni di cloud storage. Un collega usa zip() con expand_kwargs() per abbinare ogni file alla sua destinazione nel cloud:

@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):
    # Simula il caricamento di un file sul 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))

Quante istanze di upload creerà Airflow?

Questo esercizio fa parte del corso

Creare data pipeline con Airflow

Visualizza corso

esercizio interattivo pratico

Trasforma la teoria in pratica con uno dei nostri esercizi interattivi

Inizia esercizio