Détection d'anomalies en temps réel
Un pipeline de streaming permettant de surveiller les interactions sur un site e-commerce et d’identifier en temps réel les comportements suspects. Une alerte est déclenchée lorsqu’un utilisateur dépasse le seuil de 50 clics par minute, afin de détecter notamment les bots, les scrapers ...
pipeline
Architecture du système de détection d'anomalies
1. Ingestion des données — Kafka
Chaque clic effectué sur le site e-commerce est représenté par un événement contenant plusieurs informations : user_id, event, product_id et timestamp.
Ces événements sont envoyés vers un topic Kafka, qui assure leur collecte et leur transmission en temps réel. Un producteur Python simule le trafic normal des utilisateurs, tandis qu’un second simulateur génère 10 000 clics rapides afin de reproduire un comportement anormal et tester le système.
2. Traitement en temps réel — Flink
Les événements sont ensuite consommés et traités en temps réel par Apache Flink via PyFlink.
Le système utilise l’event time et des watermarks pour gérer les événements pouvant arriver avec du retard. Les clics sont regroupés dans des fenêtres temporelles d’une minute par utilisateur (tumbling windows).
Pour réduire l’utilisation de la mémoire lors du comptage des clics, le système utilise une structure Count-Min Sketch.
3. Détection et génération des alertes
À la fermeture de chaque fenêtre d’une minute, Flink vérifie le nombre de clics générés par chaque utilisateur.
Si un utilisateur dépasse le seuil de 50 clics par minute, son activité est considérée comme anormale et une alerte est générée automatiquement.
L’alerte contient notamment :
- user_id : identifiant de l’utilisateur
- window_range : période concernée
- click_count : nombre de clics détectés