DAG et les scheduler,tasks Airflow et les orchestrer,IO avec Airflow.Pour ce TP, utilisez la branch suivante :
git checkout 4_starting_orchestration
Sur cette branche, il y a maintenant :
dags/train.py qui permet d'entraîner un modèledags/predict.py qui est incomplet et qui permettra de réaliser des prédictionsformation_mlops_2/ ont été décorées avec des read et des write pour donner des fonctions function_name_with_ioLe dossier scripts contient des scripts d'entraînement et de prédiction pour notre cas d'usage de Machine Learning.
Nous allons désormais voir comment orchestrer ces tâches grâce à Airflow.
Revue de code avec les formateurs pour introduire les concepts de DAGs et de tâches dans le code.
Il n'est pas conseillé de partager en mémoire de la donnée d'une tâche à l'autre dans un DAG Airflow, il convient plutôt de les écrires dans des fichiers.
Pour répondre à ce problème, nous avons décoré la fonction de prédiction avec
Nous avons créé une fonction train_with_io & predict_with_io qui soient utilisables par le DAG Airflow.
Les prédictions réalisées sont écrites dans 2 fichiers identiques :
%Y%m%d-%H%M%SLes méthodes _with_io ont également des import lazy, c'est-à-dire qu'ils se font au runtime, plutôt qu'au chargement du script pour accélérer le dag-orchestrator.
Airflow se configure à travers un fichier de configuration situé dans le dossier airflow_home. Il est possible de configurer Airflow à travers des variables d'environnement, mais pour ce TP, nous allons utiliser le fichier de configuration.
Pour voir l'ensemble des configurations possibles, allez voir la documentation officielle.
Nous devons apporter quelques modifications au fichier de configuration actuel :
/airflow/airflow.cfg, avec l'éditeur nano : nano /airflow/airflow.cfg, ou bien avec vscode.dags_folder pour pointer sur /home/jovyan/work/Formation-MLOps-2/dags, cela permet d'indiquer à airflow où se situent vos DAGs# Fichier /airflow/airflow.cfg
[core]
# The folder where your airflow pipelines live, most likely a
# subfolder in a code repository
# This path must be absolute
dags_folder = /home/jovyan/work/Formation-MLOps-2/dags
...
# Whether to load the examples that ship with Airflow.
load_examples = False
...
Commençons par démarrer l'interface graphique d'Airflow, nous l'avons intégré dans l'environnement de TP.
Dans le Launcher, lancer le service Airflow.
L'interface graphique d'Airflow devrait s'ouvrir dans un nouvel onglet. Si le lancement indique could not start airflow in time cela peut vouloir dire qu'Airflow n'a pas encore démarré, essayer de refresh quelques secondes plus tard, sinon sollicitez votre formateur.
Les identifiants de connection à airflow sont admin:admin
L'interface vous indique que les différents services ne sont pas accessibles pour l'instant, c'est normal.
Naviguez, dans l'onglet Dags. Vous ne voyez pour l'instant pas de DAG, il faut alors lancer le dag processor, il se charge de parcourir votre dossier de dags, et de les parser.
uv run airflow dag-processor
Il ne faudrat pas fermer ce terminal, au risque d'arrêter le service.
La mise à jour des dags sera faite par ce service, qui les refresh par défaut toutes les 30 secondes. Pour forcer un refresh, vous pourrez l'arrêter et le relancer.
Nous allons maintenant lancer le scheduler, dans un nouveau terminal :
uv run airflow scheduler
Finalement, il faut lancer l'exécution : uv run airflow api-server --apps execution pour qu'un service s'occupe de réaliser les tâches.
L'interface graphique devrait désormais afficher 3 DAGs :

Afin de s'entraîner, il va nous falloir des données d'entraînement !
Elles ne sont pas versionnées dans ce repo. Télécharger les données avec la commande make dataset.
Les données sont désormais disponibles dans data/la-haute-borne-data-2017-2020.csv.
Pour lancer le DAG train:
Play (à droite de chaque ligne de DAG),
Inspecter le DAG train en cliquant sur celui-ci, la tâche prepare_features devrait avoir commencé :

Vous pouvez explorer les différentes informations, visuels qu'offre cette vue de DAGs.
Compléter le DAG dags/predict pour intégrer la fonction predict_with_io dans un opérateur, avec les bons arguments.
Lancer le dag data_denerator pour qu'il produise toutes les 2 minutes un petit jeu de données sur lequel nous pourrons faire des inférences.
Puis lancer le dag predict pour qu'il face les prédictions.
Après avoir manipulé des DAGs opérées par de la logique d'orchestration et le temps, nous vous proposons de découvrir le fonctionnement d'un système événementiel..
Pour illustrer ce à quoi ressemble une architecture événementiel, nous allons ouvrir deux terminaux.
touch /tmp/event.txtecho "Hello World" >> /tmp/event.txttail -f -n 1 /tmp/event.txt | xargs -I {} sh -c 'echo "{}" | wc -c'Cette implémentation basique est une illustration du comportement d'un système événementiel, si il y a un nouveau message il agit, sinon il ne fait rien. A la différnece d'un CRON qui tentera toujours de faire quelque chose, dont parfois constater la différence avec la précédente exé&cution.
Les systèmes événementiels tels que Kafka, RabbitMQ... offre bien entendu plus de fonctionnalité, de robustesse, de scalabilité.
Les instructions du tp suivant sont ici