Zacznij terazZacznij za darmo

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

Zobacz kurs

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()
Edytuj i uruchom kod