Quando os produtores enviam mensagens mais rápido do que os consumidores as processam, as mensagens não processadas se acumulam no broker do ApsaraMQ for RocketMQ. Esse acúmulo aumenta diretamente a latência das mensagens: o intervalo entre a produção e o consumo. O planejamento adequado de capacidade e os ajustes no consumo evitam que o acúmulo afete seus negócios.
Este tópico aborda o mecanismo de consumo do cliente TCP (Java SDK), as causas raiz do acúmulo e como prevenir e resolver esse problema.
Quando o acúmulo é crítico
Dois cenários exigem atenção especial:
Acúmulo persistente sem autorrecuperação. O sistema downstream consome mensagens em um ritmo inferior ao da produção pelo sistema upstream, e essa diferença não se reduz naturalmente.
Requisitos estritos de tempo real. Atrasos mínimos no consumo são inaceitáveis para o negócio.
Como o cliente consome mensagens
Um cliente TCP do ApsaraMQ for RocketMQ em modo push processa mensagens em duas fases:

Fase 1: Buscar mensagens do broker
O cliente usa long polling para buscar mensagens em lotes e as armazena em cache em uma fila de buffer local. Essa fase é rápida: mesmo um servidor com especificações modestas (4 vCPUs, 8 GB de memória), com uma única thread e uma única partição, atinge dezenas de milhares de transações por segundo (TPS). Com múltiplas partições, o throughput escala para centenas de milhares de TPS. Gargalos de acúmulo não ocorrem nesta fase.
Fase 2: Processar mensagens na thread de consumo
O cliente envia as mensagens em cache à thread de consumo, que executa a lógica de negócios. A velocidade de processamento depende de dois fatores:
Duração do consumo: tempo necessário para processar uma única mensagem.
Concorrência de consumo: quantidade de mensagens processadas em paralelo.
Quando a fila de buffer local enche, o cliente para de buscar novas mensagens no broker. O gargalo reside sempre na Fase 2, determinado pela duração e pela concorrência do consumo. Entre os dois fatores, a duração do consumo tem o maior impacto.
Causas raiz do acúmulo
O acúmulo de mensagens decorre de processamento lento por mensagem (duração do consumo) ou paralelismo insuficiente (concorrência de consumo). Identifique o fator que representa o gargalo antes de fazer ajustes.
|
Causa raiz |
Sintomas |
Correção principal |
|
Pico na duração do consumo |
A latência aumenta subitamente; a contagem de threads é adequada |
Investigue as dependências downstream |
|
Concorrência insuficiente |
A utilização da CPU está baixa; threads ociosas aguardando mensagens |
Aumente as threads por nó ou adicione nós |
Duração do consumo
A computação interna da CPU geralmente é insignificante comparada às operações de I/O externas. As operações mais comuns limitadas por I/O na lógica de consumo incluem:
Leituras e gravações em banco de dados (por exemplo, MySQL)
Leituras e gravações em cache (por exemplo, Redis)
Chamadas a serviços downstream (por exemplo, chamadas RPC Dubbo, chamadas de API HTTP)
Execute profiling e benchmark de cada chamada externa para estabelecer uma linha de base de desempenho. O acúmulo geralmente começa quando uma dependência downstream fica lenta ou atinge um limite de capacidade.
Exemplo. Um consumidor grava cada mensagem em um banco de dados em 1 ms sob carga normal. Durante um pico de tráfego, o banco de dados atinge seu limite de capacidade e a latência de gravação aumenta para 100 ms por mensagem — uma queda de 100 vezes na taxa de consumo. Aumentar a contagem de threads do cliente não resolve. A solução é escalar o banco de dados.
Concorrência de consumo
A concorrência determina quantas mensagens o cliente processa em paralelo:
|
Tipo de mensagem |
Concorrência efetiva |
|
Mensagens normais |
Threads por nó x Número de nós |
|
Mensagens agendadas e atrasadas |
Threads por nó x Número de nós |
|
Mensagens transacionais |
Threads por nó x Número de nós |
|
Mensagens ordenadas |
min(Threads por nó x Número de nós, Número de partições) |
A concorrência de mensagens ordenadas é limitada pelo número de partições no tópico. Para avaliar a contagem de partições, entre em contato com o Alibaba Cloud Customer Services.
Prioridade de ajuste: Escale as threads em um único nó primeiro. Adicione mais nós somente após a utilização total dos recursos de hardware do nó.
Estime a contagem ideal de threads
Para um nó com C vCPUs, em que cada mensagem consome T1 de tempo de CPU e T2 de tempo de espera de I/O:
|
Métrica |
Fórmula |
|
TPS de thread única |
1 / (T1 + T2) |
|
Contagem ideal de threads (100% de utilização da CPU) |
C x (T1 + T2) / T1 |
|
Throughput máximo do nó |
Contagem ideal de threads / (T1 + T2) |
Exemplo. Em um nó de 4 vCPUs em que T1 = 5 ms e T2 = 45 ms:
Contagem ideal de threads = 4 x (5 + 45) / 5 = 40 threads
Throughput máximo = 40 / (5 + 45) = 800 TPS
Esta fórmula assume condições ideais: sem sobrecarga de troca de threads, operações de I/O que não consomem CPU e memória suficiente. Em produção, aumente as threads gradualmente e monitore a utilização da CPU, a memória e o throughput antes de definir um valor final.
Previna o acúmulo
Execute benchmark do desempenho de consumo e planeje a capacidade antes da implantação. Uma linha de base bem estabelecida facilita a detecção de anomalias em produção.
Faça benchmark da duração do consumo
Meça a duração base do consumo por meio de testes de estresse e profiling de código. Concentre-se nos seguintes pontos:
Elimine problemas no nível de código. Verifique se há complexidade computacional excessiva, loops infinitos ou recursão profunda na lógica de consumo.
Minimize o I/O externo. Avalie se cada chamada externa (consulta ao banco de dados, busca em cache, API downstream) é estritamente necessária. Use cache local sempre que possível para eliminar chamadas redundantes.
Descarregue operações pesadas. Mova operações complexas e demoradas para processamento assíncrono. Garanta que o design assíncrono não crie condições de corrida — por exemplo, marcar o consumo como concluído antes que a operação assíncrona termine.
Defina a concorrência de consumo
Encontre a contagem ideal de threads. Aumente gradualmente as threads em um único nó enquanto monitora CPU, memória e throughput. Identifique o ponto de retornos decrescentes.
-
Calcule o número necessário de nós. Com base no throughput de nó único e no pico de tráfego upstream:
Number of nodes = Peak traffic / Single-node throughput
Lidando com o acúmulo quando ele ocorre
Passo 1: Configure alertas
Configure alertas de acúmulo de mensagens pelo recurso de monitoramento e alerta no ApsaraMQ for RocketMQ. Para obter detalhes, consulte Configurar alertas de acúmulo de mensagens.
Passo 2: Diagnostique e resolva
Quando um alerta disparar, identifique se a causa raiz é a duração do consumo ou a concorrência:
|
Diagnóstico |
Verificação |
Ação |
|
Pico na duração do consumo |
Consulte a métrica de duração do consumo. Verifique a latência do banco de dados, as taxas de acerto de cache e a integridade dos serviços downstream. |
Corrija o gargalo da dependência downstream. Consulte Consultar a duração do consumo. |
|
Concorrência de consumo insuficiente |
A utilização da CPU está baixa e as threads não estão saturadas. |
Aumente as threads por nó ou adicione mais nós. |
Para um guia completo de correção, consulte Lidar com mensagens acumuladas.