Eingaben mit zip paaren
Dein Team hat eine Pipeline, die Dateien in bestimmte Cloud-Speicherorte hochlädt. Eine Kollegin verwendet zip() mit expand_kwargs(), um jede Datei ihrem Cloud-Ziel zuzuordnen:
@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))
Wie viele Instanzen von upload wird Airflow erstellen?
Diese Übung ist Teil des Kurses
<Kurs>Data-Pipelines mit Airflow aufbauen</Kurs>Interaktive praktische Übung
Verwandle Theorie mit einer unserer interaktiven Übungen in die Praxis
Übung starten