Fler @tasks
För att fortsätta bygga ut ditt arbetsflöde behöver du lägga till ytterligare ett steg som tolkar och sparar ändringarna i den nedladdade filen. Dag:en process_sales är definierad och har redan uppgiften pull_file tillagd. Python-funktionen parse_file(inputfile, outputfile) är redan definierad åt dig.
Observera att du när du implementerar Airflow-uppgifter inte alltid behöver förstå varje enskilt steg i detalj. Så länge du vet hur du kapslar in stegen i Airflows struktur kan du bygga det arbetsflöde du behöver.
Den här övningen är en del av kursen
Introduktion till Apache Airflow i Python
Övningsinstruktioner
- Skapa en Airflow-uppgift med hjälp av metoden
parse_file. - Anropa uppgiften med de argument som behövs.
Interaktiv övning med praktiskt arbete
Testa den här övningen genom att slutföra den här exempelkoden.
@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()