L'accumulation de messages dans ApsaraMQ for Kafka survient lorsque les consommateurs ne parviennent pas à suivre le rythme des producteurs, ce qui entraîne une accumulation de messages non traités dans les partitions du topic. Cette métrique est couramment appelée retard du consommateur (consumer lag).
Sans intervention, l'accumulation de messages augmente la latence de traitement, déclenche des boucles de rééquilibrage et peut provoquer des erreurs d'épuisement de mémoire (OOM) sur les clients consommateurs.
Fonctionnement de l'accumulation de messages
Chaque partition d'un topic suit deux offsets :
Offset du consommateur : position jusqu'à laquelle un groupe de consommateurs a traité les messages.
Dernier offset : position du message produit le plus récemment.
L'écart entre ces deux valeurs représente l'accumulation pour cette partition :
Accumulation per partition = Latest offset - Consumer offset
Total accumulation = Sum of accumulation across all partitions
Une accumulation proche de zéro indique que les consommateurs suivent le rythme des producteurs. Une accumulation croissante signale que les consommateurs prennent du retard.
Exemple
Topic: test (Partition 0)
+----+----+----+----+----+----+----+
| M1 | M2 | M3 | M4 | M5 | M6 | M7 | <- 7 messages written
+----+----+----+----+----+----+----+
^ ^
Consumer offset (3) Latest offset (7)
Topic: test (Partition 1)
+----+----+----+----+----+----+
| M1 | M2 | M3 | M4 | M5 | M6 | <- 6 messages written
+----+----+----+----+----+----+
^ ^
Consumer offset (3) Latest offset (6)
Total accumulation = (7 - 3) + (6 - 3) = 7 messages
Si aucun offset de consommateur n'existe (parce que le consommateur n'a pas validé d'offset ou que l'offset a expiré) et qu'au moins un thread de consommateur du groupe est en ligne, l'accumulation est calculée comme Latest offset - Earliest offset sur toutes les partitions. Si tous les threads de consommateurs du groupe sont hors ligne, l'accumulation est de 0.
Diagnostic de la cause racine
Déterminez d'abord si les consommateurs sont actifs :
Les consommateurs sont hors ligne ou redémarrent
| Symptôme | Cause probable |
|---|---|
| Les processus consommateurs ne s'exécutent pas | Plantages, mises à jour de déploiement ou longues pauses de garbage collection (GC) |
| Les consommateurs alternent entre les états en ligne et hors ligne | Rééquilibrages fréquents causés par l'ajout ou la suppression de consommateurs dans le groupe, des délais d'expiration des battements de cœur ou une valeur faible pour session.timeout.ms |
Les consommateurs sont actifs mais prennent du retard
| Symptôme | Cause probable |
|---|---|
| Augmentation constante de l'accumulation | Débit insuffisant des consommateurs : entrées/sorties lentes, goulots d'étranglement au niveau du CPU ou de la mémoire, ou logique de traitement par message trop lente |
| Pic soudain de l'accumulation | Afflux de trafic dû à une charge de pointe ou à des importations par lots |
| Accumulation malgré des ressources adéquates | Problèmes de code tels que des boucles infinies, des exceptions non interceptées ou de longs intervalles entre les appels poll() |
| Accumulation à un rythme constant | Limitation du débit des consommateurs : le taux de consommation a atteint la limite réservée ou élastique de l'instance |
| Accumulation semblant plus importante que prévu | Validations d'offset retardées ou échouées, entraînant des extractions répétées et des chiffres de retard gonflés |
Surveillance des métriques d'accumulation
Consultez les métriques d'accumulation selon votre type d'instance :
| Type d'instance | Outil de surveillance |
|---|---|
| Abonnement ou paiement à l'utilisation | Surveillance Prometheus |
| Serverless | Tableau de bord |
Impact d'une accumulation non résolue
Une accumulation non résolue affecte le système global de plusieurs manières :
Latence accrue : le traitement retardé affecte les services en aval et les décisions commerciales.
Threads bloqués et délais d'expiration : les consommateurs saturés peuvent se bloquer, provoquant des délais d'expiration des requêtes et déclenchant des disjoncteurs.
Boucles de rééquilibrage : un traitement lent entraîne des délais d'expiration des battements de cœur, ce qui déclenche un rééquilibrage des partitions. Le rééquilibrage suspend la consommation, augmente les extractions répétées et aggrave le retard, créant ainsi une boucle de rétroaction négative.
Erreurs OOM : si un consommateur appelle
poll()mais ne traite pas les messages assez rapidement, les messages non traités s'accumulent dans le tampon mémoire du client et peuvent provoquer un débordement de tas.
Résolution de l'accumulation de messages
Augmenter le débit des consommateurs
Ajouter des instances de consommateurs : ajoutez davantage de consommateurs au même groupe. Chaque partition est attribuée à au plus un consommateur au sein d'un groupe, donc le nombre de consommateurs ne doit pas dépasser le nombre de partitions (
partitions >= consumers).Augmenter le nombre de partitions : davantage de partitions permettent un parallélisme accru. Chaque nouvelle partition peut être attribuée à un consommateur distinct.
Traiter les messages de manière asynchrone : déchargez les opérations chronophages (écritures dans la base de données, appels API) vers un pool de threads ou une file d'attente de tâches afin que
poll()retourne rapidement.Utiliser le traitement par lots : traitez plusieurs messages par itération de boucle au lieu de les traiter un par un.
Ajuster les paramètres des consommateurs
| Paramètre | Valeur par défaut | Plage recommandée | Objectif |
|---|---|---|---|
max.poll.records |
500 | 1 – 500 | Contrôle le nombre d'enregistrements renvoyés par chaque appel poll(). Des valeurs plus faibles réduisent le temps de traitement par interrogation. |
fetch.min.bytes |
1 B | 1 KB – 1 MB | Définit la quantité minimale de données que le broker collecte avant de répondre à une requête d'extraction. Des valeurs plus élevées réduisent les extractions vides et améliorent le débit. |
fetch.max.wait.ms |
500 ms | 500 ms | Définit le temps maximal pendant lequel le broker attend pour accumuler fetch.min.bytes avant de répondre. |
session.timeout.ms |
10 s | 30 s | Temps maximal pendant lequel le broker attend sans battement de cœur avant de déclarer un consommateur mort. Une valeur plus élevée empêche les expulsions injustifiées. |
heartbeat.interval.ms |
3 s | <= session.timeout.ms / 3 |
Contrôle la fréquence des battements de cœur. Doit être inférieure à session.timeout.ms pour maintenir l'appartenance au groupe. |
enable.auto.commit |
false | true | Active les validations automatiques d'offset. Empêche les échecs de validation d'offset de provoquer une fausse accumulation. |
Mesures d'urgence
Si l'accumulation est trop importante pour être résorbée par un traitement normal, réinitialisez l'offset du consommateur à la dernière position. Cette opération ignore tous les messages accumulés et reprend la consommation à partir de l'offset le plus récent.
La réinitialisation de l'offset du consommateur à la dernière position supprime tous les messages accumulés. N'utilisez cette approche que lorsque les messages accumulés ne sont plus nécessaires.
Pour effacer les alertes basées sur l'accumulation, utilisez ApsaraMQ for Kafka pour réinitialiser l'offset du consommateur d'une partition de topic à 0. Lorsque l'offset du consommateur est 0, l'accumulation signalée est 0.