Sintomas
Um alerta de acúmulo de mensagens em uma instância do Message Queue for RocketMQ indica um problema potencial. Ao acessar o console do Message Queue for RocketMQ, observe os seguintes sintomas:
Na página Group Details, o valor de Real-time Accumulated Messages do Group ID está acima do esperado.
No painel de navegação, selecione Message Tracing, clique em Create Query Task e selecione Query by Message ID. Após insira as informações necessárias, verifique se algumas mensagens foram enviadas ao broker, mas não entregues aos consumidores downstream.
Possíveis causas
No Message Queue for RocketMQ, após o envio das mensagens a um broker, o cliente configurado com um Group ID busca as mensagens no broker com base no offset atual do consumidor e as processa localmente. O acúmulo geralmente não ocorre durante a busca de mensagens pelo cliente. Esse cenário costuma resultar de capacidade insuficiente de processamento no cliente, causada por fatores como tempo de consumo prolongado ou baixa concorrência de consumidores. Para obter mais informações sobre o mecanismo de consumo e as causas do acúmulo de mensagens, consulte Problemas de acúmulo e latência de mensagens.
Solução
Se ocorrer acúmulo de mensagens, siga estas etapas para resolver o problema.
-
Identifique se o acúmulo de mensagens ocorre no servidor do Message Queue for RocketMQ ou no cliente.
Verifique o arquivo de log local do cliente
ons.loge pesquise a seguinte mensagem:the cached message count exceeds the thresholdSe essa mensagem de log aparecer, a fila de buffer local do cliente estará cheia e as mensagens estarão acumulando no lado do cliente. Nesse caso, prossiga para a Etapa 2.
Se essa mensagem de log não for encontrada, o acúmulo não estará ocorrendo no cliente. Nessa situação excepcional, entre em contato com o suporte técnico da Alibaba Cloud.
-
Verifique se o tempo de consumo das mensagens é adequado.
Se o tempo de consumo for excessivamente longo, analise o stack trace do cliente para solucionar problemas na lógica de negócio. Nesse caso, prossiga para a Etapa 3.
Se o tempo de consumo estiver normal, o acúmulo pode decorrer de concorrência insuficiente nos consumidores. Para corrigir isso, aumente gradualmente o número de threads de consumo ou escale horizontalmente os nós consumidores.
Verifique o tempo de consumo das mensagens das seguintes formas:
Acesse o console do Message Queue for RocketMQ para visualizar o rastreamento da mensagem. Na seção Consumer, visualize o consumption time de uma única mensagem. Para mais detalhes, consulte Consultar rastreamentos de mensagens. Para obter o tempo de consumo de uma mensagem específica, consulte o rastreamento e verifique o Message Processing Time nos resultados de entrega ao consumidor, na página de detalhes do rastreamento.
Acesse o console do Message Queue for RocketMQ para verificar o status do consumidor. Nas informações de conexão do cliente, confira o Response Time para obter o tempo médio de consumo. Para mais informações, consulte Visualizar status do consumidor. Na seção Consumption Statistics da janela pop-up Java Client Real-time Data, visualize métricas de cada tópico, como Response Time (ms/message), Successful Messages (messages/s), Failed Messages (messages/s) e Message Accumulation. Se o valor de Response Time estiver muito alto (por exemplo, 5003,94 ms/message) e o acúmulo de mensagens aumentar continuamente, haverá um problema de latência de consumo.
Utilize outros produtos de monitoramento, como o Application Real-Time Monitoring Service (ARMS), para instrumentar sua aplicação e coletar dados sobre o tempo de consumo de mensagens.
-
Analise o stack trace do cliente. Concentre-se apenas nas threads nomeadas ConsumeMessageThread, pois elas contêm a lógica de negócio para o consumo de mensagens. Consulte a documentação oficial do Java para determinar o estado da thread e ajustar a lógica de negócio conforme necessário.
Obtenha o stack trace do cliente das seguintes maneiras:
Acesse o console do Message Queue for RocketMQ e visualize o status do consumidor. Nas informações de conexão do cliente, localize a opção View Stack Information. Para mais detalhes, consulte Visualizar status do consumidor.
-
Use a ferramenta jstack para imprimir o stack trace.
Consulte Visualizar status do consumidor para encontrar o endereço IP do host que executa a instância consumidora com acúmulo de mensagens. Em seguida, acesse esse host.
-
Execute um dos comandos abaixo para visualizar e registrar o ID do processo (PID) do processo Java.
ps -ef |grep javajps -lm -
Execute o comando a seguir para visualizar o stack trace.
jstack -l pid > /tmp/pid.jstack -
Execute o comando abaixo para visualizar informações sobre as threads
ConsumeMessageThread.cat /tmp/pid.jstack|grep ConsumeMessageThread -A 10 --color
Os exemplos a seguir mostram stack traces anormais comuns:
-
Exemplo 1: Stack trace ocioso, sem acúmulo.
Quando um consumidor está ocioso, suas threads de consumo ficam no estado
WAITING, aguardando a recuperação de mensagens da fila de tarefas de consumo.Por exemplo, o stack trace da thread
ConsumeMessageThread_7exibeTID: 53 STATE: WAITING, e a cadeia de chamadas éThreadPoolExecutor.runWorker→ThreadPoolExecutor.getTask→LinkedBlockingQueue.take→LockSupport.park→sun.misc.Unsafe.park. -
Exemplo 2: A lógica de consumo envolve contenção de locks ou estados de espera (sleep).
A thread do consumidor fica bloqueada por uma operação interna de sleep ou wait, o que torna o consumo mais lento.
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) -
Exemplo 3: A lógica de consumo trava em uma operação com um sistema de armazenamento externo, como um banco de dados.
A thread do consumidor fica bloqueada por uma chamada HTTP externa, o que desacelera o consumo.
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)
Se o acúmulo de mensagens afetar seus serviços e as mensagens puderem ser ignoradas com segurança, utilize o recurso Redefinir offsets do consumidor. Esse recurso permite ignorar as mensagens acumuladas e retomar o consumo a partir do offset mais recente, restaurando rapidamente seus serviços. Para mais informações, consulte Redefinir offsets do consumidor.