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 Airflow planifié : 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.
Cet exercice fait partie du cours
<cours>Introduction au data engineering</cours>Instructions de l’exercice
- Complétez la tâche
etl()en utilisant les fonctions définies dans l'énoncé de l'exercice. - Configurez la dépendance correcte. Notez que la tâche
etldoit attendre la fin dewait_for_table. - La dernière ligne illustre une exécution :
etl.function()appelle la fonction Python sous-jacente à 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()