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
esercizio interattivo pratico
Trasforma la teoria in pratica con uno dei nostri esercizi interattivi
Inizia esercizio