ÎncepețiÎncepe gratuit

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

Vezi cursul

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()
Editează și rulează codul