Kom igångKom igång gratis

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

Visa kurs

Interaktiv övning med praktiskt arbete

Gör teori till handling med en av våra interaktiva övningar

Starta övningen