Ещё один @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()