O acúmulo de mensagens no ApsaraMQ for Kafka ocorre quando os consumidores não acompanham os produtores e as mensagens não processadas se acumulam nas partições do tópico. Essa métrica é conhecida como consumer lag.
Sem resolução, o acúmulo de mensagens aumenta a latência de processamento, dispara ciclos de rebalanceamento e pode causar erros de falta de memória (OOM) nos clientes consumidores.
Como funciona o acúmulo de mensagens
Cada partição em um tópico rastreia dois offsets:
Consumer offset: posição até a qual um grupo de consumidores processou as mensagens.
Latest offset: posição da mensagem produzida mais recentemente.
A diferença entre esses dois valores representa o acúmulo daquela partição:
Accumulation per partition = Latest offset - Consumer offset
Total accumulation = Sum of accumulation across all partitions
Um acúmulo próximo de zero indica que os consumidores acompanham os produtores. Um acúmulo crescente sinaliza atraso no consumo.
Exemplo
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
Se não existir um consumer offset (porque o consumidor não confirmou um offset ou porque ele expirou) e pelo menos uma thread consumidora no grupo estiver online, o sistema calculará o acúmulo como Latest offset - Earliest offset em todas as partições. Caso todas as threads consumidoras do grupo estejam offline, o acúmulo será 0.
Diagnosticar a causa raiz
Verifique primeiro se os consumidores estão em execução:
Consumidores offline ou reiniciando
|
Sintoma |
Causa provável |
|
Processos consumidores fora de execução |
Falhas, atualizações de implantação ou pausas longas de coleta de lixo (GC) |
|
Alternância entre estados online e offline |
Rebalanceamentos frequentes causados por entrada ou saída de consumidores do grupo, timeouts de heartbeat ou valor baixo de |
Consumidores em execução, mas atrasados
|
Sintoma |
Causa provável |
|
Aumento constante no acúmulo |
Throughput insuficiente dos consumidores: I/O lento, gargalos de CPU ou memória, ou lógica de processamento por mensagem muito lenta |
|
Pico repentino no acúmulo |
Surto de tráfego devido a carga máxima ou importações em lote |
|
Acúmulo mesmo com recursos adequados |
Problemas no código, como loops infinitos, exceções não capturadas ou intervalos longos entre chamadas de |
|
Acúmulo a uma taxa constante |
Limitação de taxa do consumidor: a taxa de consumo atingiu o limite reservado ou elástico da instância |
|
Acúmulo maior que o esperado |
Confirmações de offset atrasadas ou com falha, causando pulls repetidos e números de lag inflados |
Monitorar métricas de acúmulo
Visualize as métricas de acúmulo conforme o tipo da sua instância:
|
Tipo de instância |
Ferramenta de monitoramento |
|
Assinatura ou pagamento conforme o uso |
|
|
Serverless |
Impacto do acúmulo não resolvido
O acúmulo não resolvido afeta todo o sistema de várias formas:
Latência aumentada: o atraso no processamento impacta serviços downstream e decisões de negócios.
Threads bloqueadas e timeouts: consumidores sobrecarregados podem bloquear, causar timeouts de requisição e disparar circuit breakers.
Ciclos de rebalanceamento: o processamento lento leva a timeouts de heartbeat e dispara o rebalanceamento de partições. Esse processo pausa o consumo, aumenta os pulls repetidos e piora o lag, criando um ciclo de feedback negativo.
Erros de OOM: se um consumidor chama
poll()mas não processa as mensagens com rapidez suficiente, as mensagens não processadas se acumulam no buffer de memória do cliente e podem causar estouro de heap.
Resolver o acúmulo de mensagens
Dimensionar o throughput dos consumidores
Adicione instâncias consumidoras: inclua mais consumidores no mesmo grupo. Cada partição é atribuída a no máximo um consumidor dentro de um grupo; portanto, o número de consumidores não deve exceder o número de partições (
partitions >= consumers).Aumente as partições: mais partições permitem maior paralelismo. Cada nova partição pode ser atribuída a um consumidor separado.
Processe mensagens de forma assíncrona: descarregue operações demoradas (gravações em banco de dados, chamadas de API) para um pool de threads ou fila de tarefas para garantir que
poll()retorne rapidamente.Utilize processamento em lote: processe múltiplas mensagens por iteração de loop em vez de uma por vez.
Ajustar parâmetros do consumidor
|
Parâmetro |
Padrão |
Faixa recomendada |
Finalidade |
|
|
500 |
1 -- 500 |
Controla quantos registros cada chamada de |
|
|
1 B |
1 KB -- 1 MB |
Define a quantidade mínima de dados que o broker coleta antes de responder a uma requisição de fetch. Valores maiores reduzem fetches vazios e melhoram o throughput. |
|
|
500 ms |
500 ms |
Estabelece o tempo máximo que o broker aguarda para acumular |
|
|
10 s |
30 s |
Tempo máximo que o broker aguarda sem um heartbeat antes de declarar o consumidor como inativo. Um valor maior evita remoções falsas. |
|
|
3 s |
<= |
Determina a frequência do heartbeat. Deve ser menor que |
|
|
false |
true |
Ative confirmações automáticas de offset. Evita que falhas na confirmação de offset causem acúmulos falsos. |
Medidas de emergência
Se o acúmulo for grande demais para ser drenado pelo processamento normal, redefina o consumer offset para a posição mais recente. Isso ignora todas as mensagens acumuladas e retoma o consumo a partir do offset mais recente.
Redefinir o consumer offset para a posição mais recente descarta todas as mensagens acumuladas. Utilize essa abordagem apenas quando as mensagens acumuladas não forem mais necessárias.
Para limpar alertas baseados em acúmulo, use o ApsaraMQ for Kafka para redefinir o consumer offset de uma partição do tópico para 0. Quando o consumer offset é 0, o acúmulo reportado é 0.