Symptômes
Une alerte signalant une accumulation de messages sur une instance Message Queue for RocketMQ indique un problème potentiel. Après vous être connecté à la console Message Queue for RocketMQ, vous pouvez constater les symptômes suivants :
Sur la page Group Details, la valeur de Real-time Accumulated Messages pour l'ID de groupe est supérieure aux attentes.
Dans le volet de navigation, sélectionnez Message Tracing, cliquez sur Create Query Task, puis sélectionnez Query by Message ID. Après avoir saisi les informations requises, vous constatez que certains messages sont envoyés au broker mais ne sont pas distribués aux consommateurs en aval.
Causes possibles
Dans Message Queue for RocketMQ, après l'envoi des messages à un broker, un client configuré avec un ID de groupe extrait les messages du broker en fonction de l'offset de consommation actuel et les consomme localement. L'accumulation de messages ne se produit généralement pas lors de l'extraction des messages depuis le broker. Elle est habituellement due à une capacité de traitement insuffisante du client, provoquée par des facteurs tels qu'un temps de consommation élevé ou une faible concurrence des consommateurs. Pour plus d'informations sur le mécanisme de consommation et les causes de l'accumulation de messages, reportez-vous à la rubrique Problèmes d'accumulation et de latence des messages.
Solution
En cas d'accumulation de messages, suivez les étapes ci-dessous pour résoudre le problème.
-
Déterminez si l'accumulation de messages se produit sur le serveur Message Queue for RocketMQ ou sur le client.
Consultez le fichier journal local du client
ons.loget recherchez le message suivant :the cached message count exceeds the thresholdSi ce message de journal apparaît, la file d'attente tampon locale du client est pleine et les messages s'accumulent sur le client. Dans ce cas, passez à l'étape 2.
Si ce message de journal est introuvable, l'accumulation de messages ne se produit pas sur le client. Dans ce cas exceptionnel, contactez le support technique d'Alibaba Cloud.
-
Vérifiez si le temps de consommation des messages est raisonnable.
Si le temps de consommation est trop long, examinez la trace de pile du client pour déboguer votre logique métier. Dans ce cas, passez à l'étape 3.
Si le temps de consommation est normal, l'accumulation de messages peut être due à une concurrence insuffisante des consommateurs. Pour résoudre ce problème, augmentez progressivement le nombre de threads de consommation ou ajoutez des nœuds de consommateurs.
Plusieurs méthodes permettent de vérifier le temps de consommation des messages :
Connectez-vous à la console Message Queue for RocketMQ pour afficher la trace des messages. Dans la section Consumer, consultez le consumption time d'un message unique. Pour plus d'informations, reportez-vous à la rubrique Interroger les traces de messages. Pour obtenir le temps de consommation d'un message unique, interrogez la trace du message et affichez le Message Processing Time dans les résultats de distribution du consommateur sur la page des détails de la trace.
Connectez-vous à la console Message Queue for RocketMQ pour afficher l'état du consommateur. Dans les informations de connexion du client, vérifiez le Response Time pour obtenir le temps de consommation moyen. Pour plus d'informations, reportez-vous à la rubrique Afficher l'état du consommateur. Dans la section Consumption Statistics de la fenêtre contextuelle Java Client Real-time Data, vous pouvez afficher les métriques pour chaque topic, telles que Response Time (ms/message), Successful Messages (messages/s), Failed Messages (messages/s) et Message Accumulation. Si la valeur Response Time est trop élevée (par exemple, 5003,94 ms/message) et que l'accumulation de messages augmente continuellement, un problème de latence de consommation existe.
Utilisez d'autres produits de surveillance, tels que Application Real-Time Monitoring Service (ARMS), pour instrumenter votre application et collecter des données sur le temps de consommation des messages.
-
Affichez la trace de pile du client. Concentrez-vous uniquement sur les threads nommés ConsumeMessageThread. Ces threads contiennent la logique métier de consommation des messages. Consultez la documentation officielle Java pour déterminer l'état du thread et modifier votre logique métier en conséquence.
Plusieurs méthodes permettent d'obtenir la trace de pile du client :
Connectez-vous à la console Message Queue for RocketMQ et affichez l'état du consommateur. Dans les informations de connexion du client, trouvez l'option View Stack Information. Pour plus d'informations, reportez-vous à la rubrique Afficher l'état du consommateur.
-
Utilisez l'outil jstack pour imprimer la trace de pile.
Reportez-vous à la rubrique Afficher l'état du consommateur pour trouver l'adresse IP de l'hôte exécutant l'instance de consommateur présentant une accumulation de messages. Connectez-vous ensuite à cet hôte.
-
Exécutez l'une des commandes suivantes pour afficher et enregistrer l'ID de processus (PID) du processus Java.
ps -ef |grep javajps -lm -
Exécutez la commande suivante pour afficher la trace de pile.
jstack -l pid > /tmp/pid.jstack -
Exécutez la commande suivante pour afficher les informations sur les threads
ConsumeMessageThread.cat /tmp/pid.jstack|grep ConsumeMessageThread -A 10 --color
Les exemples suivants illustrent des traces de pile anormales courantes :
-
Exemple 1 : Une trace de pile inactive sans accumulation.
Lorsqu'un consommateur est inactif, ses threads de consommation sont dans l'état
WAITINGet attendent de récupérer des messages depuis la file d'attente des tâches de consommation.Par exemple, la trace de pile du thread
ConsumeMessageThread_7afficheTID: 53 STATE: WAITING, et la chaîne d'appels estThreadPoolExecutor.runWorker→ThreadPoolExecutor.getTask→LinkedBlockingQueue.take→LockSupport.park→sun.misc.Unsafe.park. -
Exemple 2 : La logique de consommation implique une contention de verrou ou des états de mise en veille.
Le thread de consommateur est bloqué par une opération interne de mise en veille ou d'attente, ce qui ralentit la consommation.
Thread: ConsumeMessageThread_16 Stack: TID: 51 STATE: TIMED_WAITING java.lang.Thread.sleep(Native Method) mqtest.DelayTest$1.consume(DelayTest.java:51) com.aliyun.openservices.ons.api.impl.rocketmq.ConsumerImpl$MessageListenerImpl.consumeMessage(ConsumerImpl.java:101) com.aliyun.openservices.shade.com.alibaba.rocketmq.client.impl.consumer.ConsumeMessageConcurrentlyService$ConsumeRequest.run(ConsumeMessageConcurrentlyService.java:415) java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) java.util.concurrent.FutureTask.run(FutureTask.java:266) java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) java.lang.Thread.run(Thread.java:748) -
Exemple 3 : La logique de consommation est bloquée sur une opération avec un système de stockage externe, tel qu'une base de données.
Le thread de consommateur est bloqué par un appel HTTP externe, ce qui ralentit la consommation.
ConsumeMessageThread_3 TID: 54 STATE: RUNNABLE java.lang.ClassLoader.loadClass(ClassLoader.java:404) sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:349) java.lang.ClassLoader.loadClass(ClassLoader.java:357) org.apache.http.impl.conn.PoolingHttpClientConnectionManager.<init>(PoolingHttpClientConnectionManager.java:174) org.apache.http.impl.conn.PoolingHttpClientConnectionManager.<init>(PoolingHttpClientConnectionManager.java:158) org.apache.http.impl.conn.PoolingHttpClientConnectionManager.<init>(PoolingHttpClientConnectionManager.java:149) org.apache.http.impl.conn.PoolingHttpClientConnectionManager.<init>(PoolingHttpClientConnectionManager.java:125) refactor.base.Tools.getHttpsClient(Tools.java:138) refactor.base.Tools.httpsPost(Tools.java:257) mqtest.DelayTest$1.consume(DelayTest.java:58) com.aliyun.openservices.ons.api.impl.rocketmq.ConsumerImpl$MessageListenerImpl.consumeMessage(ConsumerImpl.java:101) com.aliyun.openservices.shade.com.alibaba.rocketmq.client.impl.consumer.ConsumeMessageConcurrentlyService$ConsumeRequest.run(ConsumeMessageConcurrentlyService.java:415) java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) java.util.concurrent.FutureTask.run(FutureTask.java:266) java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) java.lang.Thread.run(Thread.java:748)
Si l'accumulation de messages affecte vos services et que les messages peuvent être ignorés en toute sécurité, utilisez la fonctionnalité Réinitialiser les offsets des consommateurs. Cette fonctionnalité vous permet d'ignorer les messages accumulés et de reprendre la consommation à partir du dernier offset afin de restaurer rapidement vos services. Pour plus d'informations, reportez-vous à la rubrique Réinitialiser les offsets des consommateurs.