CommencezCommencez gratuitement

Définir un DAG

Dans les exercices précédents, vous avez appliqué les trois étapes du processus ETL :

  • Extract : Extraire la table PostgreSQL film vers pandas.
  • Transform : Scinder la colonne rental_rate du DataFrame film.
  • Load : Charger le DataFrame film dans 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

Voir le cours

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 etl doit attendre la fin de wait_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()
Modifier et exécuter le code