Bắt đầu ngayBắt đầu miễn phí

Ghép cặp đầu vào với zip

Nhóm của bạn có một pipeline để tải các tệp lên các vị trí lưu trữ đám mây cụ thể. Một đồng nghiệp dùng zip() với expand_kwargs() để ghép mỗi tệp với đích lưu trữ tương ứng:

@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):
    # Giả lập việc tải tệp lên lưu trữ đám mây
    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))

Airflow sẽ tạo bao nhiêu instance của upload?

Bài tập này là một phần của khóa học

Xây dựng Data Pipeline với Airflow

Xem khóa học

Bài tập tương tác thực hành

Biến lý thuyết thành hành động với một trong các bài tập tương tác của chúng tôi

Bắt đầu bài tập