Sintomas
Quando um cliente consumidor usa a estratégia de atribuição de partições StickyAssignor, várias threads de consumidor passam a consumir a mesma partição. Esse comportamento causa processamento duplicado ou fora de ordem das mensagens.
Causa
Esse é um bug conhecido (KAFKA-7026 / KIP-341) em versões do cliente Apache Kafka anteriores à 2,3. O StickyAssignor não elimina duplicidades na atribuição de partições quando um consumidor reingressa no grupo com dados de atribuição desatualizados.
O cenário a seguir reproduz o problema:
O consumidor C1 entra em um grupo de consumidores como líder e recebe a partição
test-0.O consumidor C2 entra no mesmo grupo. O C1 mantém a partição
test-0; o C2 não recebe nenhuma partição.O C1 deixa de responder (por exemplo, devido a uma pausa longa de GC). O C2 assume a liderança e toma posse da partição
test-0.O C1 se recupera e reingressa no grupo com sua atribuição desatualizada (
test-0). Durante o rebalanceamento, tanto o C1 quanto o C2 informamtest-0como sua atribuição existente.Como o
StickyAssignornão verifica duplicatas, ele atribui a partiçãotest-0a ambos os consumidores.
Solução
Opção 1: Atualizar o cliente Kafka para a versão 2,3 ou posterior (recomendado)
A correção desse bug está disponível no Apache Kafka 2,3. Atualize a dependência do cliente Kafka na sua aplicação:
<!-- Maven example -->
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.3.0</version> <!-- or later -->
</dependency>
Opção 2: Adotar outra estratégia de atribuição de partições
Se não for possível atualizar imediatamente, mude para uma estratégia diferente de atribuição de partições. A tabela a seguir descreve as estratégias disponíveis:
|
Estratégia |
Nome da classe |
Descrição |
Compensações |
|
Range (padrão) |
|
Distribui as partições de cada tópico uniformemente entre os consumidores. |
Simples e previsível. Pode gerar distribuição desigual quando a quantidade de partições não for múltipla da quantidade de consumidores. |
|
Round-robin |
|
Atribui partições uma a uma, em esquema de rodízio, entre todos os consumidores. |
Mais equilibrada que a Range. Pode provocar mais movimentações de partições durante os rebalanceamentos. |
|
Cooperative sticky |
|
Usa a mesma lógica de balanceamento do |
Minimiza a movimentação de partições e evita o bug de atribuição duplicada. Exige uma migração em duas etapas a partir de assignors com protocolo eager. |
Para alterar a estratégia, defina a propriedade de configuração do consumidor partition.assignment.strategy:
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
"org.apache.kafka.clients.consumer.RoundRobinAssignor");
Migrar para o CooperativeStickyAssignor
Para migrar do StickyAssignor para o CooperativeStickyAssignor em um grupo de consumidores em execução sem tempo de inatividade, execute uma reinicialização gradual em duas etapas:
-
Adicione o
CooperativeStickyAssignorcomo estratégia secundária, mantendo a estratégia atual. Em seguida, realize uma reinicialização gradual de todos os consumidores.props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.StickyAssignor," + "org.apache.kafka.clients.consumer.CooperativeStickyAssignor"); -
Após todos os consumidores aplicarem a nova configuração, alterne exclusivamente para o
CooperativeStickyAssignor. Depois, faça outra reinicialização gradual.props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
Não use oStickyAssignorem versões do cliente anteriores à 2,3. Mesmo após a correção, oCooperativeStickyAssignorgeralmente é a melhor opção, pois permite rebalanceamento incremental sem pausar todos os consumidores.