Vandeplanque Antoine
Ce TP vise à construire un système de données intégrant un Data Lake fichier, un Data Warehouse MySQL, et des Consommateurs Kafka en Python. Il fait suite au TP KsqlDB 2 qui définissait les transformations initiales. Les technologies clés utilisées sont Kafka, Python (kafka-python, mysql-connector, schedule), MySQL, et un système de fichiers local pour le Data Lake.
Le flux de données est le suivant :
- Un Producer Python envoie des transactions vers Kafka (
transaction_log). - ksqlDB (défini au TP2) traite ces données et génère des topics intermédiaires et des topics "changelog" pour ses Tables.
- Des Consommateurs Python lisent les topics Kafka pertinents (source et changelogs).
- Les données lues sont stockées dans le Data Lake (
./data_lake/feed_name/date=YYYY-MM-DD/). - Les données issues des topics "changelog" des Tables ksqlDB sont synchronisées (UPSERT/DELETE) dans le Data Warehouse MySQL.
- Un Scheduler Python exécute périodiquement des tâches de maintenance (nettoyage DL).
- Data Lake (
./data_lake/) : Stockage fichier partitionné par feed et date (YYYY-MM-DD). Format JSON par message. Mode Append. - Data Warehouse (MySQL) : Base
my_dwhcontenant les tables :dwh_user_transaction_totals: Agrégats par utilisateur/type (PK: user_id, transaction_type).dwh_windowed_transaction_totals: Agrégats fenêtrés par type (PK: transaction_type, window_start).dwh_datalake_permissions: Gestion des accès au DL.
- Producer Kafka (
kafka_producer_transaction.py) : Script Python générant des messages Kafka. Volume augmenté en modifiant le nombre de messages et/ou le délai d'envoi. - Consommateurs Kafka (
kafka_consumers.py) : Script Python utilisantkafka-python. Lance des threads pour lire des topics spécifiques définis dansTOPIC_PROCESSING_CONFIG. Écrit au DL et/ou effectue des UPSERT/DELETE dans le DWH MySQL via des fonctions réutilisables. Gère les tombstones pour les DELETEs DWH. - Gouvernance & Sécurité :
- Suppression DL : Script périodique (
scheduler_cleanup.py) supprimant les partitions datées viafindou équivalent. - Permissions DL : Table
dwh_datalake_permissionsdans MySQL pour contrôler l'accès futur. - Ajout Feed : Mise à jour de
TOPIC_PROCESSING_CONFIGdans le consommateur, et création/configuration éventuelle de la table DWH cible.
- Suppression DL : Script périodique (
- Orchestration (
scheduler_cleanup.py) : Utilise la bibliothèqueschedulepour lancer la fonction de nettoyage du Data Lake toutes les 10 minutes.
- Python 3.7+
- Docker & Docker Compose (pour Kafka/ksqlDB)
- Serveur MySQL
- Dépendances Python :
pip install kafka-python mysql-connector-python schedule python-dotenv requests
- Démarrez Kafka, ksqlDB, MySQL.
- Lancez les consommateurs :
python kafka_consumers.py - Lancez le scheduler :
python scheduler_cleanup.py - Lancez le producer :
python kafka_producer_transaction.py