НачатьНачать бесплатно

Ещё один @task

Чтобы продолжить реализацию рабочего процесса, нужно добавить ещё один шаг — разбор и сохранение изменений из загруженного файла. DAG process_sales уже определён, и задача pull_file в него добавлена. Функция Python также уже написана за вас: parse_file(inputfile, outputfile).

Обратите внимание: при реализации задач Airflow вам не всегда нужно разбираться в деталях каждого отдельного шага. Достаточно понимать, как обернуть эти шаги в структуру Airflow, — и вы сможете реализовать нужный рабочий процесс.

Это упражнение является частью курса

Введение в Apache Airflow на Python

Посмотреть курс

Инструкции к упражнению

  • Создайте задачу Airflow, используя метод parse_file.
  • Вызовите задачу с необходимыми аргументами.

Интерактивное практическое упражнение

Попробуйте выполнить это упражнение, дополнив этот пример кода.

@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()
Редактировать и запускать код