Plus de @tasks
Pour poursuivre la mise en place de votre flux de travail, vous devez ajouter une autre étape pour analyser et enregistrer les changements du fichier téléchargé. Le Dag process_sales est défini et la tâche pull_file y est déjà ajoutée. Dans ce cas, la fonction Python est déjà définie pour vous : parse_file(inputfile, outputfile).
Notez que, souvent, lorsque vous implémentez des tâches Airflow, vous ne comprendrez pas nécessairement chacune des étapes qui vous sont fournies. Tant que vous savez encapsuler ces étapes dans la structure d'Airflow, vous pourrez mettre en œuvre le flux de travail souhaité.
Cette activité fait partie du cours
Introduction à Apache Airflow en Python
Instructions de l’exercice
- Créez une tâche Airflow en utilisant la méthode
parse_file. - Appelez la tâche avec les arguments nécessaires.
Exercice interactif pratique
Essayez cet exercice en complétant ce code d’exemple.
@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()