Více @tasks
Aby ses mohl/a posunout dál v implementaci svého workflow, potřebuješ přidat další krok pro zpracování a uložení změn staženého souboru. DAG process_sales je definovaný a má již přidaný task pull_file. Pythonová funkce parse_file(inputfile, outputfile) je pro tebe připravena předem.
Měj na paměti, že při implementaci Airflow tasků nemusíš vždy rozumět každému jednotlivému kroku. Stačí vědět, jak jednotlivé kroky zabalit do struktury Airflow — a zvládneš sestavit jakékoli požadované workflow.
Toto cvičení je součástí kurzu
Úvod do Apache Airflow v Pythonu
Pokyny k cvičení
- Vytvoř Airflow task pomocí metody
parse_file. - Zavolej task s potřebnými argumenty.
Interaktivní cvičení na vyzkoušení si v praxi
Vyzkoušejte si toto cvičení dokončením tohoto ukázkového kódu.
@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()