Začněte nyníZačněte zdarma

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

Zobrazit kurz

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()
Upravit a spustit kód