Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Empty file modified Changelog.md
100644 → 100755
Empty file.
Empty file modified README.md
100644 → 100755
Empty file.
Empty file added data_ingestion/__init__.py
Empty file.
88 changes: 88 additions & 0 deletions data_ingestion/ingest_data.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
import requests
import pandas as pd
from loguru import logger

BASE_URL = "http://127.0.0.1:8001"

def get_users():
"""Fetches all users from the API and returns them as a pandas DataFrame."""
all_users = []
page = 1
while True:
response = requests.get(f"{BASE_URL}/users?page={page}&size=100")
if response.status_code == 200:
data = response.json()
users = data["items"]
if not users:
break
all_users.extend(users)
page += 1
logger.info(f"Fetched page {page} of users, total users fetched: {len(all_users)}")
else:
logger.warning(f"Failed to fetch users: {response.status_code}")
logger.error(f"Error details: {response.text}")
break
return pd.DataFrame(all_users)


def get_tracks():
"""Fetches all tracks from the API and returns them as a pandas DataFrame."""
all_tracks = []
page = 1
while True:
response = requests.get(f"{BASE_URL}/tracks?page={page}&size=100")
if response.status_code == 200:
data = response.json()
tracks = data["items"]
if not tracks:
break
all_tracks.extend(tracks)
page += 1
logger.info(f"Fetched page {page} of tracks, total tracks fetched: {len(all_tracks)}")
else:
logger.warning(f"Failed to fetch tracks: {response.status_code}")
logger.error(f"Error details: {response.text}")
break
return pd.DataFrame(all_tracks)

def get_listen_history():
"""Fetches all listen history from the API and returns it as a pandas DataFrame."""
all_listen_history = []
page = 1
while True:
response = requests.get(f"{BASE_URL}/listen_history?page={page}&size=100")
if response.status_code == 200:
data = response.json()
listen_history = data["items"]
if not listen_history:
break
all_listen_history.extend(listen_history)
page += 1
logger.info(f"Fetched page {page} of listen history, total records fetched: {len(all_listen_history)}")
else:
logger.warning(f"Failed to fetch listen history: {response.status_code}")
logger.error(f"Error details: {response.text}")
break
return pd.DataFrame(all_listen_history)

def main():
"""init log file """
logger.add("data_ingestion.log", rotation="1 MB", level="INFO")
logger.info("Starting data ingestion process")
logger.info("##########################################################################")
"""Main function to ingest data from all endpoints."""
users_df = get_users()
tracks_df = get_tracks()
listen_history_df = get_listen_history()

print("Users DataFrame:")
print(users_df.head())
print("\nTracks DataFrame:")
print(tracks_df.head())
print("\nListen History DataFrame:")
print(listen_history_df.head())
logger.info("Data ingestion process completed successfully")
logger.info("##########################################################################")

if __name__ == "__main__":
main()
143 changes: 139 additions & 4 deletions docs/ANSWERS.md
100644 → 100755
Original file line number Diff line number Diff line change
@@ -1,23 +1,158 @@
# Réponses du test

## _Utilisation de la solution (étape 1 à 3)_
etape 1:
utilisation de uv en installant une version 3.9 puis activation de cette environement : source .env_uv_3919/bin/activate
installation des paquets uv pip install -r requirements.txt

etape 2:
lancement du serveur:
une erreur dans les methodes getters du main utilise la classe Page qui n'avait pas les methodes get_tracks, get_users et get_listen_history
modification du code en utilisant la classe CustomPage a la place de Page
le port 8000 etant occupe,
j'ai redirige en local le serveur sur le port 8001 avec python -m uvicorn main:app --port 8001

Le serveur demarre et l'url http://127.0.0.1:8001/docs affiche l'API

etape 3 :
je teste trois aspects :
- la reponse 200 du serveur
- la reponse en terme de data
- le schema de la data retourne

_Inscrire la documentation technique_

## Questions (étapes 4 à 7)

### Étape 4

_votre réponse ici_
nous partons de la data d'origine et du besoin de faire des recommandations a partir des chansons ecoutes.
une chanson seraecoute a une date d par un ou plusieurs users. Un user ecoutera une ou plusieurs chansons. En modele relationnel , nous aurions une table pivot qui correspondrait au coeur de ce qui va faire l'objet des calculs pour les recommandations : la table user_history.
Dans votre cas, la table d'association pourrait être `user_listen_history` (ou `ecoutes`):

**Tables principales :**

* **`users`**
* `user_id` (PK)
* `user_name`
* ...

* **`songs`**
* `song_id` (PK)
* `title`
* `artist`
* ...

**Table d'association :**

* **`user_listen_history`**
* `listen_id` (PK, optionnel, mais bonne pratique pour une clé primaire unique)
* `user_id` (FK vers `users.user_id`)
* `song_id` (FK vers `songs.song_id`)
* `listen_timestamp` (Quand l'écoute a eu lieu - très important pour les recommandations basées sur l'historique récent)
* `play_duration_seconds` (Durée d'écoute, si pertinent)
* ...

Avec cette structure, vous pouvez facilement savoir quelles chansons un utilisateur a écoutées, et quels utilisateurs ont écouté une chanson donnée, ainsi que des détails sur chaque événement d'écoute.

Ainsi ce modele peut etre porte par un SGDBR comme Postgresql.

Les donnees massives issue de l'API , elles seraint stocker dans un data lake tel AWS S3.
Quand aux donnees issue des calculs de machine learning seraient stocker dans une base de donnees cle valeur comme MongoDB:
"Calcul des recommandations terminé et servi via MongoDB"
```

**Exemple de document dans MongoDB :**

```json
{
"_id": "user_12345", // La clé est le user_id pour un lookup O(1)
"recommendations": [
"song_id_abc",
"song_id_def",
"song_id_ghi",
// ...
],
"model_version": "v1.2.3",
"calculation_date": "2025-08-04T22:00:00Z"
}
```
ainsi , on a un workflow hybride Data lake de donnees massives sur un cloud + les donnees maitres sur une base relationnelles + les resultats des calculs de recommandations fait a partir de la table user_history_listen

### Étape 5

_votre réponse ici_
le systeme de surveillance de la sante du pipel;ine va reposer sur :
1- un systeme de logging , etape par etape , de chaque composant du pipeline
2- un dashboard via un outil BI pour visualiser les metriques ce qui permet d'avoir une big picture du pipeline
3- des alertes envoyeant aux intervenants cles lies aux actions cles du pipelines dans un canala specifique de communication(email, tchat ...)

venons en aux metriques:
- les metriques lie au pipeline , aux ressources materielles consommes, aux statutss desou ko des etapes du pipeline.

- les metriques lie a la data : le nombre de lignes , la preence de null , la coherence du schema , les types de donnees



### Étape 6

_votre réponse ici_
Pour automatiser le calcul des recommandations de type scoring , je mettrais en place un workflow orchestré qui se déclenche après la réussite du pipeline d'ingestion de données.

**Architecture et Flux de Travail :**

```
[FIN du Pipeline d'Ingestion] -> [DÉCLENCHEMENT du Workflow de Calcul]
|
V
[1. Tâche de Calcul Batch (ex: Spark,..)]
| a. Charge le dernier modèle validé
| b. Charge les nouvelles données d'historique des utilisateurs.
| c. Pour chaque utilisateur, calcule les N recommandations.
| d. Écrit les résultats dans une table "staging".
|
V
[2. Tâche de Validation des Résultats]
| a. Vérifie la qualité des recommandations (ex: nombre d'utilisateurs avec des recommandations, pas de valeurs nulles).
|
V
[3. Tâche de Keep or Replace]
| a. Si la validation est réussie, remplace l'ancienne table de recommandations par la nouvelle (staging).
| b. Archive l'ancienne table pour analyse.
|
V
[4. Notification] -> [Canal Slack/Email: "Calcul des recommandations terminé avec succès"]

### Étape 7

_votre réponse ici_
la le reentrainement du modele de recommandation va necessite d'evaluer le modele existant et le modele issue de donnees recentes. Ces evaluations vontse faire sur les metriques de performances des modeles selon un jeu de donnees test.

[DÉCLENCHEUR (ex: Hebdomadaire, ou baisse de performance)]
|
V
[1. Tâche d'Extraction et Préparation des Données]
| a. Crée un jeu de données d'entraînement (ex: 90 derniers jours) et de test (ex: 7 derniers jours).
|
V
[2. Tâche d'Entraînement Parallèle]
| a. Entraîne un nouveau "modèle candidat" sur le jeu de données d'entraînement.
|
V
[3. Tâche d'Évaluation Comparative]
| a. Charge le "modèle en production" actuel.
| b. Évalue le "modèle candidat" ET le "modèle en production" sur le même jeu de données de test.
| c. Compare leurs métriques de performance.
|
V
[4. Tâche de Validation et Enregistrement (Conditionnelle)]
| a. SI le candidat est meilleur que la production (selon un seuil défini, ex: +5% de précision):
| i. Enregistre le "modèle candidat".
| ii. Attribue au nouveau modèle le tag "staging" ou "validation".
| b. SINON:
| i. Garde le modèle en production et alerte l'équipe (le modèle n'apprend plus).
|
V
[5. Tâche de Déploiement/Promotion (Manuelle ou Automatique)]
| a. Un Data Scientist valide les métriques du modèle en "staging".
| b. Promotion du modèle : le tag passe de "staging" à "production".
|
V
[6. Notification] -> [Rapport de réentraînement envoyé sur Slack/Email]
77 changes: 77 additions & 0 deletions notebooks/01_users_ingestion.ipynb
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
{
"cells": [
{
"cell_type": "code",
"execution_count": 1,
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
" id first_name last_name email gender \\\n",
"0 8231 Gabrielle Mcguire melinda03@example.org Male \n",
"1 30893 Morgan Hill jshaw@example.com Bigender \n",
"2 13477 Tiffany Wiley crosslisa@example.com Gender questioning \n",
"3 24153 Jonathan Kerr terryjessica@example.com Genderqueer \n",
"4 16236 Robert Hensley dennisoliver@example.net Female \n",
"\n",
" favorite_genres created_at updated_at \n",
"0 Hip Hop 2025-01-05T05:13:07.505433 2025-07-30T15:54:19.063299 \n",
"1 Jazz 2023-10-17T19:55:22.351170 2025-04-28T09:18:47.655815 \n",
"2 Country 2023-12-11T06:51:59.356696 2025-06-28T22:13:55.643900 \n",
"3 Indie 2024-12-04T01:05:42.812122 2024-12-27T23:50:52.591169 \n",
"4 Blues 2024-08-14T02:08:12.199497 2025-01-02T19:13:52.554275 \n"
]
}
],
"source": [
"import requests\n",
"import pandas as pd\n",
"\n",
"BASE_URL = \"http://127.0.0.1:8001\"\n",
"\n",
"def get_users():\n",
" all_users = []\n",
" page = 1\n",
" while True:\n",
" response = requests.get(f\"{BASE_URL}/users?page={page}&size=100\")\n",
" if response.status_code == 200:\n",
" data = response.json()\n",
" users = data[\"items\"]\n",
" if not users:\n",
" break\n",
" all_users.extend(users)\n",
" page += 1\n",
" else:\n",
" print(f\"Failed to fetch users: {response.status_code}\")\n",
" break\n",
" return pd.DataFrame(all_users)\n",
"\n",
"users_df = get_users()\n",
"print(users_df.head())"
]
}
],
"metadata": {
"kernelspec": {
"display_name": ".env_uv_3919",
"language": "python",
"name": "python3"
},
"language_info": {
"codemirror_mode": {
"name": "ipython",
"version": 3
},
"file_extension": ".py",
"mimetype": "text/x-python",
"name": "python",
"nbconvert_exporter": "python",
"pygments_lexer": "ipython3",
"version": "3.9.19"
}
},
"nbformat": 4,
"nbformat_minor": 4
}
Loading