Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Best practices for producers

Última atualização: Jun 27, 2026

Ao enviar grandes volumes de mensagens com um produtor Kafka, configurações inadequadas de tentativas, lotes ou confirmações podem causar perda de dados, latência excessiva ou erros de memória insuficiente. Este tópico explica como configurar o produtor do ApsaraMQ for Kafka para garantir envios confiáveis e throughput ideal. Todos os exemplos utilizam o cliente Java. Os conceitos fundamentais aplicam-se a outras linguagens, mas os nomes dos parâmetros e os detalhes de implementação podem variar.

Enviar uma mensagem

Toda interação do produtor começa com producer.send(), que aceita um ProducerRecord contendo tópico, partição, timestamp, chave e valor.

Future<RecordMetadata> metadataFuture = producer.send(new ProducerRecord<String, String>(
        topic,   // The message topic.
        null,   // The partition number. Set this to null to let the producer assign a partition.
        System.currentTimeMillis(),   // The timestamp.
        String.valueOf(value.hashCode()),   // The message key.
        value   // The message value.
));

O método send() é assíncrono. Para bloquear a execução até obter o resultado, chame:

RecordMetadata metadata = metadataFuture.get(timeout, TimeUnit.MILLISECONDS);

Para exemplos completos do SDK, consulte Visão geral do SDK.

Chave e valor

As mensagens no ApsaraMQ for Kafka 0.10.2 têm dois campos:

  • Chave -- Identificador da mensagem. Defina uma chave exclusiva por mensagem para rastrear seu ciclo de vida nos logs de envio e consumo.

  • Valor -- Conteúdo da mensagem.

Em cenários de alto volume sem necessidade de roteamento por chave, omita esse campo e utilize a estratégia de particionamento sticky para melhorar a eficiência do agrupamento em lotes.

Importante

O ApsaraMQ for Kafka 0.11.0 e versões posteriores suportam headers. Para usar headers, atualize o servidor para a versão 2.2.0.

Segurança de threads

A instância do produtor é thread-safe e pode enviar mensagens para qualquer tópico. Use apenas um produtor por aplicação.

Configure tentativas

Em ambientes distribuídos, problemas de rede podem interromper o envio de mensagens. A falha ocorre em dois momentos possíveis: a mensagem chega ao broker, mas a confirmação (ACK) se perde; ou a mensagem nunca alcança o broker.

O ApsaraMQ for Kafka usa uma arquitetura de rede com endereço IP virtual (VIP) que encerra automaticamente conexões ociosas. Clientes inativos frequentemente recebem o erro connection reset by peer. As tentativas de reenvio tratam essas falhas transitórias.

Parâmetro

Descrição

Valor recomendado

retries

Número máximo de tentativas para um envio com falha.

Use o padrão (definido pela versão do cliente).

retry.backoff.ms

Atraso entre as tentativas, em milissegundos.

1000

Confirme

O parâmetro acks controla quantas réplicas devem confirmar uma gravação antes da resposta do broker. Escolha a configuração conforme seus requisitos de durabilidade e desempenho.

Configuração

Comportamento

Throughput

Durabilidade

acks=0

Nenhuma resposta do broker necessária.

Mais alto

Mais baixa -- perda de dados provável em qualquer falha

acks=1

Resposta após o nó primário gravar os dados.

Médio

Média -- perda de dados se o nó primário falhar antes da replicação

acks=all

Resposta após o nó primário e todos os nós de réplica sincronizada (ISR) gravarem os dados.

Mais baixo

Mais alta -- perda de dados apenas se os nós primário e réplica falharem simultaneamente

Para melhorar o desempenho de envio, defina acks=1.

Otimizar o agrupamento em lotes

O produtor agrupa mensagens destinadas à mesma partição em lotes antes do envio. Lotes maiores reduzem requisições de rede, diminuem o uso de CPU e melhoram o throughput e a latência. Já os lotes pequenos provocam filas de requisições tanto no cliente quanto no servidor.

Dois parâmetros controlam o agrupamento em lotes:

Parâmetro

Descrição

Padrão

Recomendação

batch.size

Tamanho máximo do lote por partição, em bytes. Uma requisição de rede é acionada quando o lote atinge esse tamanho.

16384 (16 KB)

Mantenha o padrão de 16384. Valores muito baixos degradam o desempenho e a estabilidade.

linger.ms

Tempo máximo que uma mensagem aguarda no buffer antes do envio, em milissegundos. Ao expirar esse tempo, o produtor envia o lote independentemente do batch.size.

0

Defina entre 100 e 1000.

O envio do lote ocorre quando qualquer um dos limiares é atingido, prevalecendo o que ocorrer primeiro. Para equilibrar throughput e latência, defina batch.size=16384 e linger.ms=1000.

Use um cliente da versão 2,4 ou posterior, que ativa a estratégia de particionamento sticky por padrão para reduzir ainda mais os envios fragmentados.

Estratégia de particionamento sticky

Somente mensagens enviadas para a mesma partição são agrupadas em um lote; portanto, a estratégia de particionamento afeta diretamente a eficiência desse agrupamento.

Mensagens com chave

Para mensagens com chave, o produtor calcula o hash da chave e selecione a partição com base nesse resultado. Mensagens com a mesma chave sempre vão para a mesma partição.

Mensagens sem chave

Antes do Kafka 2,4, a estratégia padrão para mensagens sem chave era round-robin: cada mensagem ia para a próxima partição em sequência. Isso dispersava as mensagens por todas as partições, gerando muitos lotes pequenos e aumentando a latência.

O Kafka 2,4 introduziu a estratégia de particionamento sticky (KIP-480) para resolver esse problema. Em vez de alternar partições a cada mensagem, o produtor mantém-se em uma única partição até que o lote atual esteja cheio e então selecione aleatoriamente uma nova partição. Com o tempo, as mensagens continuam distribuídas uniformemente por todas as partições, mas os lotes tornam-se muito maiores. Essa abordagem evita desbalanceamento de mensagens entre partições, reduz a latência e melhora o desempenho geral do serviço.

Ative o particionamento sticky

  • Cliente versão 2,4 ou posterior: O particionamento sticky é o padrão. Nenhuma configuração necessária.

  • Cliente anterior à versão 2,4: Implemente um particionador personalizado e configure-o via partitioner.class. O exemplo abaixo alterna partições em um intervalo de tempo configurável:

public class MyStickyPartitioner implements Partitioner {

    // Records the time of the last partition switch.
    private long lastPartitionChangeTimeMillis = 0L;
    // Records the current partition.
    private int currentPartition = -1;
    // The partition switch interval. Set the interval as needed.
    private long partitionChangeTimeGap = 100L;

    public void configure(Map<String, ?> configs) {}

    /**
     * Compute the partition for the given record.
     *
     * @param topic The topic name
     * @param key The key to partition on (or null if no key)
     * @param keyBytes serialized key to partition on (or null if no key)
     * @param value The value to partition on or null
     * @param valueBytes serialized value to partition on or null
     * @param cluster The current cluster metadata
     */
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {

        // Get all partition information.
        List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
        int numPartitions = partitions.size();

        if (keyBytes == null) {
            List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic);
            int availablePartitionSize = availablePartitions.size();

            // Check the current active partitions.
            if (availablePartitionSize > 0) {
                handlePartitionChange(availablePartitionSize);
                return availablePartitions.get(currentPartition).partition();
            } else {
                handlePartitionChange(numPartitions);
                return currentPartition;
            }
        } else {
            // For a message with a key, select a partition based on the key's hash value.
            return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
        }
    }

    private void handlePartitionChange(int partitionNum) {
        long currentTimeMillis = System.currentTimeMillis();

        // If the time since the last switch exceeds the switch interval, switch to the next partition. Otherwise, use the current partition.
        if (currentTimeMillis - lastPartitionChangeTimeMillis >= partitionChangeTimeGap
            || currentPartition < 0 || currentPartition >= partitionNum) {
            lastPartitionChangeTimeMillis = currentTimeMillis;
            currentPartition = Utils.toPositive(ThreadLocalRandom.current().nextInt()) % partitionNum;
        }
    }

    public void close() {}

}

Evitar erros de memória insuficiente (OOM)

O produtor armazena mensagens em cache na memória antes de enviar os lotes. Se o cache crescer excessivamente, ocorre um erro de memória insuficiente (OOM).

Parâmetro

Descrição

Padrão

Recomendação

buffer.memory

Memória total disponível para o buffer de envio do produtor, em bytes. Se o buffer for muito pequeno, a alocação de memória pode demorar, afetando o desempenho de envio e causando timeouts.

33554432 (32 MB)

Defina pelo menos como batch.size x número de partições x 2. O padrão de 32 MB é suficiente para um único produtor.

Importante

Execute múltiplos produtores na mesma Java Virtual Machine (JVM) multiplica o uso de memória. Cada produtor aloca sua própria buffer.memory -- quatro produtores com o padrão de 32 MB consomem 128 MB de heap. Em produção, geralmente basta um único produtor por aplicação. Caso seja necessário executar vários produtores, reduza a buffer.memory de cada um ou aumente o tamanho do heap da JVM adequadamente.

Ordenação de partições

Dentro de uma única partição, as mensagens são armazenadas e consumidas na ordem de envio.

Por padrão, o ApsaraMQ for Kafka não garante ordenação estrita dentro de uma partição. Durante uma atualização ou failover, um pequeno número de mensagens pode ficar desordenado ao ser redirecionado para outra partição. Esse compromisso melhora a disponibilidade.

Para impor ordenação estrita dentro de uma partição, selecione local storage ao criar o tópico.