Więcej dekoratorów @task
Aby kontynuować budowanie swojego przepływu pracy, musisz dodać kolejny krok – parsowanie i zapisywanie zmian z pobranego pliku. DAG process_sales jest już zdefiniowany i zawiera zadanie pull_file. W tym przypadku funkcja Pythona jest już dla ciebie gotowa: parse_file(inputfile, outputfile).
Pamiętaj, że podczas implementowania zadań w Airflow nie zawsze musisz rozumieć każdy szczegół poszczególnych kroków. Wystarczy, że wiesz, jak opakować je w strukturę Airflow – i już możesz z powodzeniem zbudować dowolny przepływ pracy.
To ćwiczenie jest częścią kursu
Wprowadzenie do Apache Airflow w Pythonie
Instrukcje do ćwiczenia
- Utwórz zadanie Airflow, używając metody
parse_file. - Wywołaj zadanie z wymaganymi argumentami.
Interaktywne ćwiczenie praktyczne
Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.
@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()