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

Nhiều @tasks hơn

Để tiếp tục triển khai workflow, bạn cần thêm một bước nữa để phân tích và lưu các thay đổi của tệp đã tải xuống. Dag process_sales đã được định nghĩa và đã có sẵn task pull_file. Trong trường hợp này, hàm Python đã được định nghĩa sẵn cho bạn, parse_file(inputfile, outputfile).

Lưu ý rằng khi triển khai các task trong Airflow, bạn không phải lúc nào cũng hiểu rõ từng bước được giao. Chỉ cần bạn biết cách bao bọc các bước đó trong cấu trúc của Airflow, bạn sẽ có thể triển khai workflow mong muốn.

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

Giới thiệu về Apache Airflow bằng Python

Xem khóa học

Hướng dẫn bài tập

  • Tạo một task Airflow bằng phương thức parse_file.
  • Gọi task với các đối số cần thiết.

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

Hãy thử làm bài tập này bằng cách hoàn thành đoạn mã mẫu này.

@dag(dag_id='process_sales')
def process_sales():
    # Decorate parse_file as a task
    ____
    def parse_file(inputfile: str, outputfile: str):
        with open(inputfile) as infile:
            data = json.load(infile)
            with open(outputfile, 'w') as outfile:
                json.dump(data, outfile)

    pull_file('http://dataserver/sales.json', 'latestsales.json')
    # Call the parse_file task
    ____('latestsales.json', 'latestsales_parsed.json')

process_sales()
Chỉnh sửa và Chạy Mã