Définir un DAG
Dans les exercices précédents, vous avez appliqué les trois étapes du processus ETL :
- Extract : Extraire la table PostgreSQL
filmverspandas. - Transform : Scinder la colonne
rental_ratedu DataFramefilm. - Load : Charger le DataFrame
filmdans un entrepôt de données PostgreSQL.
Les fonctions extract_film_to_pandas(), transform_rental_rate() et load_dataframe_to_film() sont définies dans votre environnement de travail. Dans cet exercice, vous allez les transformer en un pipeline planifié Airflow : la tâche etl s'exécute après qu'une tâche wait_for_table a signalé que la table source est prête.
Cette activité fait partie du cours
Introduction à l'ingénierie des données
Instructions de l’exercice
- Complétez la tâche
etl()en utilisant les fonctions décrites dans l'énoncé de l'exercice. - Configurez la bonne dépendance. Notez que la tâche
etldoit attendre la fin dewait_for_table. - La dernière ligne montre une exécution d'exemple :
etl.function()appelle la fonction Python simple derrière la tâche, ce qui fait tourner le pipeline une fois ici dans la console.
Exercice interactif pratique
Essayez cet exercice en complétant ce code d’exemple.
# Define the ETL task
@task(task_id="etl_film")
def etl():
film_df = ____()
film_df = ____(____)
____(____)
# Add the task to the DAG
@dag(dag_id="etl",
start_date=datetime(2024, 1, 1),
schedule="0 0 * * *")
def etl_dag():
wait_for_table = EmptyOperator(task_id="wait_for_table")
wait_for_table >> ____
etl_dag()
# Sample run
etl.function()