ПочатиПочніть безкоштовно

Більше про @tasks

Щоб продовжити реалізацію вашого робочого процесу, додайте ще один крок для розбору та збереження змін у завантаженому файлі. 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()
Редагувати та запускати код