Більше про @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()