Este tópico descreve como configurar os parâmetros do cliente para o ApsaraMQ for Kafka. A configuração adequada dos parâmetros afeta diretamente o throughput de mensagens, a confiabilidade da entrega e a estabilidade do consumidor. As seções a seguir abordam os parâmetros de produtor e consumidor, com valores recomendados e orientações de ajuste para cargas de trabalho em produção.
Parâmetros do produtor
Entrega de mensagens
acks
Define quantos reconhecimentos do broker o produtor exige antes de considerar um envio bem-sucedido.
|
Valor |
Comportamento |
Compromisso |
|
|
Nenhum reconhecimento do broker. |
Maior throughput, maior risco de perda de dados. |
|
|
Reconhecimento após o líder gravar os dados. |
Equilíbrio entre throughput e durabilidade. Pode haver perda de dados se o líder falhar antes da replicação pelos seguidores. |
|
|
Reconhecimento após o líder e todas as réplicas sincronizadas gravarem os dados. |
Menor throughput, maior durabilidade. A perda de dados ocorre apenas se o líder e todas as réplicas sincronizadas falharem simultaneamente. |
Valor recomendado: 1 para a maioria das cargas de trabalho que priorizam throughput em vez de durabilidade estrita.
retries
Número máximo de tentativas do produtor para reenviar uma mensagem com falha. Um valor mais alto ajuda o produtor a se recuperar de falhas transitórias do broker, como eleições de líder. Combine este parâmetro com retry.backoff.ms para controlar o ritmo das novas tentativas.
retry.backoff.ms
Atraso entre as tentativas de reenvio. Um valor muito baixo pode causar tempestades de tentativas durante failovers do broker.
|
Recomendado |
Padrão |
Unidade |
|
1000 |
-- |
milissegundos |
Agrupamento e throughput
O agrupamento reduz a sobrecarga de rede ao combinar vários registros em uma única solicitação. Dois parâmetros controlam o envio do lote: tamanho e tempo.
batch.size
Tamanho máximo de um único lote por partição. Quando um lote atinge esse tamanho, o produtor o envia imediatamente.
|
Tipo |
Padrão |
Valores válidos |
Unidade |
|
int |
16384 |
[0,...] |
bytes |
Mantenha o valor padrão de 16384 para a maioria das cargas de trabalho. Um valor menor aumenta o número de solicitações de rede e reduz o throughput. Ao aumentar o batch.size, certifique-se de que o buffer.memory seja suficiente para acomodar os lotes maiores.
linger.ms
Tempo máximo que o produtor aguarda para preencher um lote antes de enviá-lo. Esse mecanismo funciona de forma semelhante ao algoritmo de Nagle no TCP: assim que um lote atinge o batch.size, o envio é imediato, independentemente do temporizador de espera. Se o lote ainda estiver abaixo do batch.size quando o linger.ms expirar, o produtor enviará o que tiver acumulado até o momento.
|
Recomendado |
Padrão |
Unidade |
|
100 a 1000 |
0 |
milissegundos |
Um valor maior para linger.ms aumenta a eficiência do agrupamento e o throughput, às custas da latência por mensagem.
Gerenciamento de memória
buffer.memory
Memória total que o produtor aloca para armazenar em buffer os registros não enviados em todas as partições. Se esse pool se esgotar, o método send() bloqueia ou lança uma exceção, dependendo do max.block.ms. Um buffer subdimensionado causa alocação lenta de memória, redução de throughput ou tempos limite de envio.
Unidade: bytes. Padrão: 33554432 (32 MB).
Fórmula de dimensionamento:
buffer.memory >= batch.size x number_of_partitions x 2
Por exemplo, com batch.size=16384 e 50 partições:
16384 x 50 x 2 = 1,638,400 bytes (~1.6 MB minimum)
Ao aumentar o batch.size para melhorar o throughput, dimensione o buffer.memory proporcionalmente.
Particionamento
partitioner.class
Determina como o produtor atribui registros às partições. A estratégia de particionamento sticky reduz o número de lotes incompletos ao preencher o lote de uma partição antes de passar para a próxima.
|
Versão do cliente Kafka |
Estratégia padrão |
|
2.4 e posteriores |
Particionamento sticky (padrão) |
|
Anterior a 2.4 |
Round-robin |
Se o cliente produtor for anterior à versão 2.4, defina explicitamente o particionador sticky para melhorar a eficiência do agrupamento.
Parâmetros do consumidor
Ajuste de busca
Estes parâmetros controlam a quantidade de dados que o consumidor recupera por solicitação de busca. O ajuste afeta tanto o throughput quanto a latência.
fetch.min.bytes
Quantidade mínima de dados que o broker acumula antes de retornar uma resposta de busca. Um valor maior reduz a frequência de buscas e a sobrecarga de CPU do broker. Isso melhora o throughput, mas aumenta a latência de mensagem ponta a ponta. Unidade: bytes.
Avalie a taxa de mensagens do produtor antes de definir este valor. Se as mensagens chegarem lentamente, um fetch.min.bytes grande adiciona atraso desnecessário.
fetch.max.wait.ms
Tempo máximo que o broker aguarda para acumular fetch.min.bytes antes de retornar uma resposta. Unidade: milissegundos.
O comportamento varia conforme o tipo de armazenamento:
Armazenamento local: O broker aguarda até que o
fetch.min.bytesseja atingido ou ofetch.max.wait.msexpire, o que ocorrer primeiro.Armazenamento em nuvem: O broker retorna uma resposta imediatamente quando novos dados chegam, independentemente do
fetch.min.bytes.
max.partition.fetch.bytes
Quantidade máxima de dados que o broker retorna por partição em uma única busca. Unidade: bytes.
Gerenciamento de sessão e rebalanceamento
Parâmetros de sessão e polling mal configurados são a causa mais comum de rebalanceamentos inesperados do consumidor. Um rebalanceamento pausa todo o consumo no grupo até a reatribuição das partições; portanto, evite acionar rebalanceamentos desnecessários.
session.timeout.ms
Tempo máximo entre heartbeats antes que o broker considere o consumidor inativo e acione um rebalanceamento.
|
Recomendado |
Intervalo válido |
Padrão |
Unidade |
|
30000 a 60000 |
6000 a 300000 |
10000 |
milissegundos |
No cliente Java 0.10.1 e posteriores, uma thread de segundo plano dedicada envia heartbeats independentemente do
poll()
. Em versões anteriores do Java ou clientes não Java, os heartbeats são enviados durante as chamadas de
poll()
, portanto, o
session.timeout.ms
deve considerar tanto o tempo de processamento de dados quanto o intervalo de heartbeat.
Dica:
Defina o
heartbeat.interval.ms
como no máximo um terço do
session.timeout.ms
. Por exemplo, se o
session.timeout.ms
for 45000, defina o
heartbeat.interval.ms
como 15000 ou menos.
max.poll.records
Número máximo de registros retornados em uma única chamada de poll(). Se o consumidor não conseguir processar essa quantidade de registros antes do prazo da próxima chamada de poll(), o broker o considerará inativo e acionará um rebalanceamento.
Fórmula de dimensionamento:
max.poll.records < messages_per_thread_per_second x consumer_threads x session_timeout_seconds
Por exemplo, com 500 msg/s por thread, 4 threads e um tempo limite de sessão de 45 segundos:
500 x 4 x 45 = 90,000
Defina o max.poll.records abaixo desse valor para garantir que o consumidor sempre termine o processamento antes que a sessão expire.
max.poll.interval.ms
Intervalo máximo entre chamadas consecutivas de poll() antes que o broker remova o consumidor do grupo. Este parâmetro aplica-se apenas ao cliente Java 0.10.1 e posteriores, onde os heartbeats são executados em uma thread separada.
|
Padrão |
Unidade |
|
300000 |
milissegundos |
Fórmula de dimensionamento:
max.poll.interval.ms > time_per_record x max.poll.records
Na maioria dos casos, o padrão de 300000 (5 minutos) é suficiente. Aumente-o apenas se a lógica de processamento for excepcionalmente lenta.
Gerenciamento de offset
enable.auto.commit
Controla se o consumidor faz commit automático dos offsets em um intervalo fixo.
|
Valor |
Comportamento |
|
|
Os offsets são confirmados automaticamente a cada |
|
|
A aplicação deve chamar |
auto.commit.interval.ms
Intervalo para commits automáticos de offset quando enable.auto.commit é true.
|
Padrão |
Unidade |
|
1000 |
milissegundos |
Um intervalo menor reduz a janela para mensagens duplicadas após uma falha, mas aumenta o número de solicitações de commit para o broker.
auto.offset.reset
Determina o que acontece quando o consumidor não tem nenhum offset confirmado ou quando o offset confirmado é inválido (por exemplo, o offset foi excluído devido a políticas de retenção).
|
Valor |
Comportamento |
|
|
Inicia o consumo a partir do offset mais recente. Apenas novas mensagens. |
|
|
Inicia o consumo a partir do offset disponível mais antigo. Reprocessa todas as mensagens retidas. |
|
|
Lança uma exceção. Use esta opção quando a aplicação gerencia offsets manualmente. |
Valor recomendado: latest. O uso de earliest faz com que o consumidor reprocesse todas as mensagens retidas sempre que encontrar um offset inválido. Isso pode levar a processamento duplicado e a um pico no atraso do consumidor.
Se a aplicação lida com o gerenciamento de offset manualmente (por exemplo, armazenando offsets em um banco de dados externo), defina este valor como none.
Referência rápida de parâmetros
Parâmetros do produtor
|
Parâmetro |
Padrão |
Recomendado |
Unidade |
|
|
-- |
|
-- |
|
|
-- |
Defina com base nos requisitos de disponibilidade |
-- |
|
|
100 |
1000 |
ms |
|
|
16384 |
16384 (padrão) |
bytes |
|
|
0 |
100 a 1000 |
ms |
|
|
33554432 |
>= |
bytes |
|
|
Sticky (2.4+) |
Particionamento sticky |
-- |
Parâmetros do consumidor
|
Parâmetro |
Padrão |
Recomendado |
Unidade |
|
|
1 |
Ajuste com base na taxa de mensagens do produtor |
bytes |
|
|
500 |
-- |
ms |
|
|
1048576 |
-- |
bytes |
|
|
10000 |
30000 a 60000 |
ms |
|
|
3000 |
<= 1/3 do |
ms |
|
|
500 |
Consulte a fórmula de dimensionamento |
-- |
|
|
300000 |
300000 (padrão) |
ms |
|
|
true |
Depende da semântica de entrega |
-- |
|
|
1000 |
1000 (padrão) |
ms |
|
|
|
|
-- |