Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Message accumulation

Última atualização: Jun 27, 2026

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
Nota

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 session.timeout.ms

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 poll()

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

Monitoramento Prometheus

Serverless

Painel

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

max.poll.records

500

1 -- 500

Controla quantos registros cada chamada de poll() retorna. Valores menores reduzem o tempo de processamento por poll.

fetch.min.bytes

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.

fetch.max.wait.ms

500 ms

500 ms

Estabelece o tempo máximo que o broker aguarda para acumular fetch.min.bytes antes de responder.

session.timeout.ms

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.

heartbeat.interval.ms

3 s

<= session.timeout.ms / 3

Determina a frequência do heartbeat. Deve ser menor que session.timeout.ms para manter a associação ao grupo.

enable.auto.commit

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.

Aviso

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.