Lorsque les producteurs envoient des messages plus rapidement que les consommateurs ne les traitent, les messages non traités s'accumulent sur le broker ApsaraMQ for RocketMQ. Ce backlog, appelé accumulation de messages, augmente directement la latence : le délai entre la production et la consommation d'un message. Une planification adéquate de la capacité et un réglage précis de la consommation évitent que cette accumulation n'impacte votre activité.
Cette rubrique présente le mécanisme de consommation du client TCP (SDK Java), les causes profondes de l'accumulation, ainsi que les méthodes pour la prévenir et la résoudre.
Quand l'accumulation devient critique
Deux scénarios nécessitent une attention particulière :
Backlog persistant sans récupération automatique. Le système aval consomme les messages plus lentement que le système amont ne les produit, et l'écart ne se résorbe pas de lui-même.
Exigences strictes en temps réel. Même de brefs retards de consommation sont inacceptables pour l'activité métier.
Fonctionnement de la consommation côté client
Un client TCP ApsaraMQ for RocketMQ en mode push traite les messages en deux phases :

Phase 1 : Récupération des messages depuis le broker
Le client utilise le long polling pour extraire les messages par lots et les met en cache dans une file d'attente locale. Cette phase est rapide : même un serveur aux spécifications modestes (4 vCPU, 8 Go de mémoire) avec un seul thread et une seule partition peut atteindre des dizaines de milliers de transactions par seconde (TPS). Avec plusieurs partitions, le débit monte à plusieurs centaines de milliers de TPS. Les goulots d'étranglement liés à l'accumulation ne surviennent pas durant cette phase.
Phase 2 : Traitement des messages dans le thread de consommation
Le client soumet les messages mis en cache au thread de consommation, qui exécute la logique métier. La vitesse de traitement dépend de deux facteurs :
Durée de consommation : temps nécessaire pour traiter un seul message.
Concurrence de consommation : nombre de messages traités en parallèle.
Lorsque la file d'attente locale est pleine, le client cesse d'extraire de nouveaux messages du broker. Le goulot d'étranglement se situe toujours dans la phase 2 et dépend de la durée de consommation et de la concurrence. Entre ces deux facteurs, la durée de consommation a l'impact le plus significatif.
Causes profondes de l'accumulation
L'accumulation de messages provient soit d'une lenteur de traitement par message (durée de consommation), soit d'un parallélisme insuffisant (concurrence de consommation). Identifiez le facteur limitant avant d'ajuster les paramètres.
| Cause profonde | Symptômes | Solution principale |
|---|---|---|
| Pic de durée de consommation | La latence augmente soudainement ; le nombre de threads est suffisant | Analyser les dépendances aval |
| Concurrence insuffisante | L'utilisation du CPU est faible ; les threads sont inactifs en attente de messages | Augmenter le nombre de threads par nœud ou ajouter des nœuds |
Durée de consommation
Le calcul interne du CPU est généralement négligeable par rapport aux E/S externes. Les opérations liées aux E/S les plus courantes dans la logique de consommation sont :
Lectures et écritures de base de données (par exemple, MySQL)
Lectures et écritures de cache (par exemple, Redis)
Appels de services aval (par exemple, appels RPC Dubbo, appels API HTTP)
Profilez et effectuez des tests de performance (benchmark) sur chaque appel externe pour établir une référence. L'accumulation commence généralement lorsqu'une dépendance aval ralentit ou atteint sa limite de capacité.
Exemple. Un consommateur écrit chaque message dans une base de données en 1 ms sous charge normale. Lors d'un pic de trafic, la base de données atteint sa limite de capacité et la latence d'écriture passe à 100 ms par message, soit une baisse de 100 fois du taux de consommation. Augmenter le nombre de threads du client n'apporte aucune amélioration. La solution consiste à augmenter la capacité de la base de données.
Concurrence de consommation
La concurrence détermine le nombre de messages que le client traite en parallèle :
| Type de message | Concurrence effective |
|---|---|
| Messages normaux | Threads par nœud x Nombre de nœuds |
| Messages planifiés et différés | Threads par nœud x Nombre de nœuds |
| Messages transactionnels | Threads par nœud x Nombre de nœuds |
| Messages ordonnés | min(Threads par nœud x Nombre de nœuds, Nombre de partitions) |
La concurrence des messages ordonnés est plafonnée par le nombre de partitions dans le topic. Pour évaluer le nombre de partitions, contactez le service client Alibaba Cloud.
Priorité de réglage : Augmentez d'abord le nombre de threads sur un seul nœud. N'ajoutez davantage de nœuds qu'une fois les ressources matérielles du nœud pleinement utilisées.
Estimer le nombre optimal de threads
Pour un nœud disposant de C vCPU, où chaque message prend T1 de temps CPU et T2 de temps d'attente E/S :
| Métrique | Formule |
|---|---|
| TPS par thread unique | 1 / (T1 + T2) |
| Nombre optimal de threads (utilisation CPU à 100 %) | C x (T1 + T2) / T1 |
| Débit maximal du nœud | Nombre optimal de threads / (T1 + T2) |
Exemple. Sur un nœud à 4 vCPU où T1 = 5 ms et T2 = 45 ms :
Nombre optimal de threads = 4 x (5 + 45) / 5 = 40 threads
Débit maximal = 40 / (5 + 45) = 800 TPS
Cette formule suppose des conditions idéales : aucun overhead de changement de contexte, des opérations E/S ne consommant pas de CPU et suffisamment de mémoire. En production, augmentez progressivement le nombre de threads et surveillez l'utilisation du CPU, la mémoire et le débit avant de fixer une valeur définitive.
Prévenir l'accumulation
Évaluez les performances de consommation et planifiez la capacité avant le déploiement. Une référence bien établie facilite la détection des anomalies en production.
Benchmarker la durée de consommation
Mesurez la durée de consommation de référence via des tests de charge et le profilage du code. Concentrez-vous sur les points suivants :
Éliminer les problèmes au niveau du code. Vérifiez l'absence de complexité computationnelle excessive, de boucles infinies ou de récursions profondes dans la logique de consommation.
Minimiser les E/S externes. Déterminez si chaque appel externe (requête de base de données, recherche dans le cache, API aval) est strictement nécessaire. Utilisez la mise en cache locale lorsque cela est possible pour éliminer les appels redondants.
Décharger les opérations lourdes. Déplacez les opérations complexes et chronophages vers un traitement asynchrone. Assurez-vous que la conception asynchrone ne crée pas de conditions de course, par exemple en marquant la consommation comme terminée avant la fin de l'opération asynchrone.
Définir la concurrence de consommation
Trouver le nombre optimal de threads. Augmentez progressivement le nombre de threads sur un seul nœud tout en surveillant le CPU, la mémoire et le débit. Identifiez le point de rendements décroissants.
-
Calculer le nombre de nœuds requis. En fonction du débit par nœud et du trafic amont de pointe :
Number of nodes = Peak traffic / Single-node throughput
Gérer l'accumulation lorsqu'elle survient
Étape 1 : Configurer les alertes
Configurez les alertes d'accumulation de messages via la fonctionnalité de surveillance et d'alerte dans ApsaraMQ for RocketMQ. Pour plus de détails, consultez la section Configurer les alertes d'accumulation de messages.
Étape 2 : Diagnostiquer et résoudre
Lorsqu'une alerte se déclenche, identifiez si la cause profonde est liée à la durée de consommation ou à la concurrence :
| Diagnostic | Vérification | Action |
|---|---|---|
| Pic de durée de consommation | Interrogez la métrique de durée de consommation. Vérifiez la latence de la base de données, les taux de hit du cache et l'état des services aval. | Corrigez le goulot d'étranglement de la dépendance aval. Consultez la section Interroger la durée de consommation. |
| Concurrence de consommation insuffisante | L'utilisation du CPU est faible et les threads ne sont pas saturés. | Augmentez le nombre de threads par nœud ou ajoutez davantage de nœuds. |
Pour un guide complet de remédiation, consultez la section Gérer les messages accumulés.