Para ihop indata med zip
Ditt team har en pipeline som laddar upp filer till specifika platser i molnlagring. En kollega använder zip() med expand_kwargs() för att para ihop varje fil med sitt molnlagringsmål:
@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))
Hur många instanser av upload kommer Airflow att skapa?
Den här övningen är en del av kursen
Bygg datapipelines med Airflow
Interaktiv övning med praktiskt arbete
Gör teori till handling med en av våra interaktiva övningar
Starta övningen