Rebalanceamentos frequentes de consumidores geralmente resultam de processamento lento de mensagens, parâmetros mal configurados ou versão desatualizada do cliente. A causa raiz e a correção variam conforme a versão do cliente.
Sintoma
Ocorrem rebalanceamentos frequentes no cliente consumidor durante o uso do ApsaraMQ for Kafka.
Causa
A causa raiz depende da versão do cliente, pois o mecanismo de heartbeat varia entre versões.
Clientes anteriores à versão 0.10.2: Não há thread dedicado para heartbeat. Os heartbeats são enviados dentro da chamada
poll(). Se o processamento da mensagem demorar mais quesession.timeout.ms, o broker não recebe heartbeat, remove o consumidor do grupo e aciona um rebalanceamento.Clientes da versão 0.10.2 ou posterior: Uma thread dedicada de heartbeat executa independentemente do processamento de mensagens. Contudo, se o consumidor não chamar
poll()dentro do intervalo definido pormax.poll.interval.ms(padrão: 5 minutos), o cliente sai do grupo de consumidores e aciona um rebalanceamento. O valor padrão demax.poll.interval.msé 5 minutos. Isso geralmente ocorre quando o processamento de um lote de mensagens leva muito tempo.
As mensagens de erro a seguir associam-se comumente a rebalanceamentos. Identificar o erro específico ajuda a localizar a causa raiz mais rapidamente:
**
NOT_COORDINATOR/join group failed**: Indica mudança no Group Coordinator (por exemplo, devido a failover do broker) ou remoção do consumidor do grupo por timeout de heartbeat. O consumidor tenta entrar novamente no grupo e aciona um rebalanceamento.**
MemberIdRequiredException/consumer poll timeout**: Em clientes anteriores à versão 0.10.2, isso geralmente resulta de um consumidor travado que bloqueia o heartbeat. Na versão 0.10.2 ou posterior, indica que o processamento lento fez o consumidor excedermax.poll.interval.mssem chamarpoll(), levando o cliente a sair do grupo de consumidores.**
CommitFailedException**: Indica que o consumidor levou tempo excessivo para processar mensagens ou falhou ao enviar heartbeats a tempo. O broker remove o consumidor do grupo (acionando um rebalanceamento) antes da conclusão do commit do offset.
Parâmetros principais
Os parâmetros a seguir controlam o comportamento de rebalanceamento:
|
Parâmetro |
Versões aplicáveis |
Descrição |
|
|
Todas |
Timeout de sessão. Se nenhum heartbeat for recebido neste período, o broker remove o consumidor do grupo. |
|
|
Todas |
Número máximo de mensagens retornadas por chamada ao método |
|
|
0.10.2 ou posterior |
Intervalo máximo entre duas chamadas consecutivas ao método |
|
|
0.10.2 ou posterior |
Intervalo de envio de heartbeats do consumidor ao broker. Deve ser menor que |
Soluções
-
Ajuste os valores dos parâmetros/@cmd
Defina os seguintes parâmetros conforme a versão do seu cliente:
**
session.timeout.ms**Versão do cliente
Valor recomendado
Anterior a 0.10.2
Maior que o tempo de processamento de um lote de mensagens, mas não superior a 30 segundos. Recomenda-se 25 segundos.
0.10.2 ou posterior
Mantenha o valor padrão de 10 segundos.
**
max.poll.records**Defina este valor bem abaixo do resultado da seguinte fórmula:
max.poll.records << messages_per_thread_per_second * number_of_threads * max.poll.interval.ms**
max.poll.interval.ms** (apenas versão 0.10.2 ou posterior)Defina este valor acima do resultado da seguinte fórmula:
max.poll.interval.ms > max.poll.records / (messages_per_thread_per_second * number_of_threads)**
heartbeat.interval.ms** (versão 0.10.2 ou posterior)Defina
heartbeat.interval.mscom no máximo 10.000 milissegundos (10 segundos). Este valor controla a frequência de envio de heartbeats ao broker e deve ser sempre menor quesession.timeout.ms. Uma prática comum é defini-lo como um terço desession.timeout.ms.Configuração de parâmetros no Spring Kafka
Se usar Spring Kafka, configure todos os parâmetros acima no arquivo
application.yml:spring: kafka: consumer: properties: session.timeout.ms: 10000 heartbeat.interval.ms: 3000 max.poll.interval.ms: 300000 max.poll.records: 50Ajuste os valores para corresponder ao seu throughput de consumo. O valor de
max.poll.interval.msdeve exceder o tempo necessário para processar um único lote de mensagens demax.poll.records. -
Melhore a velocidade de consumo e separe as threads de processamento/@cmd
Transfira o processamento de mensagens para uma thread dedicado para que a thread do consumidor chame
poll()conforme programado. Isso evita que o processamento lento bloqueie heartbeats ou exceda o intervalo de poll. -
Reduza o número de tópicos por grupo de consumidores/@cmd
Mantenha cada grupo de consumidores inscrito em no máximo cinco tópicos. Para estabilidade ideal, inscreva apenas um tópico por grupo de consumidores.
-
Atualize para a versão 0.10.2 ou posterior/@cmd
Clientes anteriores à versão 0.10.2 não possuem thread dedicada de heartbeat, tornando-os vulneráveis a timeouts durante processamento intenso. A atualização para a versão 0.10.2 ou posterior desacopla o envio de heartbeats do processamento de mensagens e elimina essa classe de rebalanceamento.
FAQ
P: Como reduzir temporariamente os rebalanceamentos para acelerar o consumo de mensagens acumuladas?
Siga as etapas abaixo para minimizar rebalanceamentos e limpar uma fila de mensagens acumuladas o mais rápido possível:
Use a versão 0.10.2 ou posterior do cliente. Clientes anteriores à versão 0.10.2 não possuem thread dedicada de heartbeat e tendem a sofrer rebalanceamentos sob carga elevada. Atualize antes de ajustar outros parâmetros.
**Aumente
max.poll.interval.ms** para um valor maior que o tempo máximo necessário para processar um único lote de mensagens. Isso evita que o consumidor seja considerado não responsivo durante processamento intenso.**Diminua
max.poll.records** para reduzir o número de mensagens buscadas por chamada depoll(). Lotes menores terminam mais rápido e reduzem o risco de excedermax.poll.interval.ms.Reduza o número de tópicos inscritos por grupo de consumidores. Mantenha as inscrições limitadas a cinco tópicos ou menos. Para estabilidade ideal, inscreva um tópico por grupo de consumidores.
Evite bloquear a thread do consumidor. Mova a lógica de processamento de mensagens para uma thread assíncrona separada para garantir que a thread do consumidor permaneça livre para chamar
poll()conforme programado.