Mai multe @tasks
Pentru a continua implementarea workflow-ului tău, trebuie să adaugi un pas suplimentar care să parseze și să salveze modificările din fișierul descărcat. DAG-ul process_sales este deja definit și include task-ul pull_file. În acest caz, funcția Python este deja scrisă pentru tine: parse_file(inputfile, outputfile).
Rețineți că, atunci când implementezi task-uri în Airflow, nu trebuie neapărat să înțelegi fiecare pas în detaliu. Atâta timp cât știi cum să înglobezi pașii în structura Airflow, vei putea implementa orice workflow dorești.
Acest exercițiu face parte din cursul
Introducere în Apache Airflow în Python
Instrucțiuni pentru exercițiu
- Creează un task Airflow folosind metoda
parse_file. - Apelează task-ul cu argumentele necesare.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
@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()