Reconstruction de la détection d'incidents avec Kafka, Flink et OpenTelemetry
Une équipe a réduit la latence de détection des incidents de 40 secondes à moins de 10 en remplaçant un agrégateur Node.js par Apache Flink sur Kubernetes et OpenTelemetry.
Traduit automatiquement depuis l'original anglais.
Une petite équipe d'ingénierie a reconstruit sa plateforme automatisée de détection d'incidents, réduisant la latence entre l'événement et la métrique de plus de 40 secondes à moins de 10. La nouvelle architecture repose sur Apache Kafka, Apache Flink exécuté sur Kubernetes et OpenTelemetry pour traiter des milliards d'événements opérationnels côté client chaque jour. Bien que le système ait amélioré la vitesse et l'isolation, l'équipe rapporte que les taux de rappel ont fluctué entre 64 % et 86 %, soulignant la complexité persistante d'une détection précise des anomalies.
Ce qui s'est passé
L'organisation exploite plus de dix produits cloud servant des millions de locataires. Auparavant, leur pile de surveillance utilisait un agrégateur Node.js alimenté par une file d'attente cloud partagée, fonctionnant sur environ 90 machines virtuelles. Ce système hérité souffrait d'une forte latence, de problèmes de « voisin bruyant » où le trafic d'autres locataires causait des retards, et de coûts qui augmentaient linéairement avec chaque nouveau produit intégré. Les coûts annuels de fonctionnement avaient grimpé de 120 000 $ à 230 000 $, et la couche cache atteignait fréquemment 100 % d'utilisation du CPU lors de modifications courantes.
Pour remédier à ces limitations, l'équipe a conçu un nouveau pipeline axé sur cinq objectifs : une latence inférieure à 10 secondes, des chemins de consommation dédiés pour l'isolation, des écritures idempotentes pour garantir la correction lors des rejeux, une évolutivité des coûts liée au volume plutôt qu'au nombre de fonctionnalités, et une exploitabilité pilotée par configuration. Le résultat est un système où l'intégration d'une nouvelle expérience ne nécessite qu'une demande de fusion (pull request) pour mettre à jour la configuration, et non un déploiement complet. L'équipe a mesuré les performances mois par mois pendant dix-huit mois, notant que si la vitesse s'est considérablement améliorée, la précision reste en dessous de leur objectif.
Comment cela fonctionne
La couche amont utilise Apache Kafka comme bus d'événements. Au lieu de consommer l'intégralité du flux massif, l'équipe a implémenté un filtre d'abonnement côté serveur défini dans le code. Ce filtre autorise uniquement des produits et expériences spécifiques tout en écartant le trafic expérimental et synthétique. Les données filtrées arrivent sur un topic Kafka dédié avec une période de rétention de sept jours, servant de fenêtre de rejeu pour le débogage ou la récupération. Deux entrées latérales alimentent le pipeline : un service de contexte de locataire fournissant des métadonnées telles que le shard et la région, et un dépôt de configuration.
Au cœur se trouve un seul job Apache Flink 1.20 déployé via le Flink Kubernetes Operator. Ce job filtre les anciens événements et les erreurs, enrichit les données avec le contexte du locataire à l'aide d'un sidecar asynchrone doté de disjoncteurs (circuit breakers), et agrège les métriques. Il utilise des schémas HyperLogLog pour estimer les utilisateurs distincts impactés sans double comptage, atteignant une marge d'erreur d'environ 1,5 %. L'état est géré dans RocksDB avec des points de contrôle (checkpoints) toutes les 30 secondes vers le stockage objet, garantissant une sémantique exactly-once pour les fichiers Parquet et des écritures idempotentes vers un magasin clé-valeur. Les métriques sont exportées via OpenTelemetry vers une base de données temporelle compatible Prometheus.
Le plan de décision, appelé AutoHOT, s'exécute comme un service Go dans deux régions. Il consomme les alertes des détecteurs, quantifie l'impact en interrogeant le magasin agrégé et applique une matrice de sévérité. Pour prévenir les faux positifs, il vérifie les pics d'activité utilisateur sur quinze minutes avant de créer un ticket. Le moteur utilise des verrous distribués pour s'assurer qu'une seule région traite une alerte à la fois, empêchant les incidents dupliqués lors des bascules (failovers). Il surveille également les silences dans la télémétrie, reconnaissant qu'une absence de données peut indiquer un shard de base de données hors ligne.
Détails clés
- La latence est passée de plus de 40 secondes à moins de 10 secondes pour le traitement événement-vers-métrique.
- Le système traite des milliards d'événements quotidiennement grâce à un seul job Apache Flink 1.20 sur Kubernetes.
- Les schémas HyperLogLog permettent des comptes d'utilisateurs distincts fusionnables à travers les locataires et les régions avec une erreur d'environ ~1,5 %.
- La configuration est gérée via un fichier YAML généré à partir de la configuration produit, permettant un chargement à chaud sans redéploiements.
- Le rappel dans le périmètre a atteint un pic de 86 % mais s'est stabilisé autour de 64 % durant les mois difficiles, indiquant une marge d'amélioration.
- Une vérification approfondie synthétique est effectuée toutes les 15 minutes pour valider l'ensemble du pipeline d'alerte, de l'injection à la clôture.
Pourquoi c'est important
Pour les ingénieurs construisant des plateformes d'observabilité, cette étude de cas démontre les compromis entre l'agrégation de type batch et le traitement de flux en temps réel. Passer à Flink a permis à l'équipe de découpler les coûts du nombre de fonctionnalités intégrées, un facteur critique pour les produits SaaS en croissance. L'utilisation d'OpenTelemetry pour une visibilité de bout en bout garantit que le système de surveillance lui-même peut être surveillé, prévenant les défaillances silencieuses qui érodent la confiance dans les tableaux de bord. Cependant, les taux de rappel fluctuants rappellent aux constructeurs que des données plus rapides ne signifient pas automatiquement une meilleure logique de détection.
Le choix architectural d'utiliser des puits (sinks) idempotents et des topics Kafka dédiés répond aux points de douleur courants des systèmes distribués : duplication de données et contention de ressources. En traitant les UID des opérateurs comme une API stable et en ajustant les limites de l'autoscaler, l'équipe a éliminé les redémarrages fréquents qui perturbaient auparavant l'état. Cette approche offre un modèle pour les équipes aux prises avec des piles de surveillance héritées, bruyantes, coûteuses et lentes, qui reposent sur des files d'attente partagées et des caches en mémoire.
Ce que vous pouvez faire
- Évaluez votre pipeline d'événements actuel pour identifier les dépendances aux files d'attente partagées pouvant causer des problèmes de voisin bruyant lors des charges maximales.
- Envisagez d'utiliser HyperLogLog ou des structures de données probabilistes similaires si vous avez besoin de comptes distincts à travers des systèmes distribués.
- Implémentez un filtrage côté serveur au niveau du bus de messages pour réduire le volume de traitement aval et les coûts.
- Concevez vos jobs de traitement de flux avec des puits idempotents pour permettre des rejeux sûrs sans double comptage des impacts.
- Ajoutez des tests synthétiques de bout en bout qui injectent de fausses alertes pour vérifier que votre pipeline de détection fonctionne même lorsque les incidents réels sont rares.
- Séparez les opérateurs dans votre graphe de flux selon leurs profils de ressources, en distinguant les transformations limitées par le CPU de l'enrichissement limité par les E/S.


