Produzindo mensagens

Um produtor é um aplicativo que publica fluxos de mensagens nos tópicos do Kafka. Essas informações se concentram na interface de programação Java que faz parte do projeto Apache Kafka. Os conceitos se aplicam a outros idiomas também, mas os nomes são, às vezes, um pouco diferentes.

Nas interfaces de programação, uma mensagem é chamada de registro. Por exemplo, a classe Java org.apache.kafka.clients.producer.ProducerRecord é usada para representar uma mensagem do ponto de vista da API producer. Os termos registro e mensagem podem ser usados de forma intercambiável, mas essencialmente um registro é usado para representar uma mensagem.

Quando um produtor se conecta ao Kafka, ele faz uma conexão de autoinicialização inicial. Essa conexão pode ser qualquer um dos servidores no cluster. O produtor solicita as informações de partição e de liderança sobre o tópico no qual deseja publicar. Em seguida, o produtor estabelece outra conexão com o líder da partição e pode publicar mensagens. Essas ações acontecem automaticamente e internamente quando seu produtor se conecta ao cluster do Kafka.

Para assegurar a disponibilidade, os brokers Kafka replicam mensagens, para que, se um broker estiver indisponível, os outros ainda possam receber mensagens de produtores e enviá-las para os consumidores Event Streams usa um fator de replicação de 3, significando que cada mensagem é armazenada em três brokers. Quando uma mensagem for enviada para o líder de partição, essa mensagem não estará imediatamente disponível para os consumidores. O líder anexa o registro para a mensagem para a partição, designando o próximo número de deslocamento para essa partição. Depois que todos os seguidores das réplicas em sincronia replicaram o registro e reconheceram que escreveram o registro em suas réplicas, o registro é confirmado e fica disponível para os consumidores.

Cada mensagem é representada como um registro que compreende duas partes: chave e valor. A chave é comumente usada para dados sobre a mensagem e o valor é o corpo da mensagem. Como muitas ferramentas do ecossistema Kafka (como conectores de outros sistemas) usam apenas o valor e ignoram a chave, é melhor colocar todos os dados da mensagem no valor e usar a chave para particionamento ou compactação do registro. Não confie em tudo que lê do Kafka para usar a chave.

Muitos outros sistemas de mensagens também têm uma forma de carregar outras informações juntamente com as mensagens. A versão 0.11 do Kafka introduz cabeçalhos de registro para essa finalidade.

Você pode achar útil ler essas informações juntamente com consumindo mensagens em Event Streams.

Definições de configuração

Muitas configurações de configuração existem para o produtor. Você pode controlar aspectos do produtor, incluindo lotes, novas tentativas e confirmação de mensagens. Aqui estão as mais importantes:

Definições de configuração do produtor
Nome Descrição Valores válidos Padrão
key.serializer A classe usada para serializar chaves. Classe Java que implementa a interface Serializer, como org.apache.kafka.common.serialization.StringSerializer. Sem padrão - você deve especificar um valor.
value.serializer A classe usada para serializar valores. Classe Java que implementa a interface Serializer, como org.apache.kafka.common.serialization.StringSerializer. Sem padrão - você deve especificar um valor.
racks O número de servidores necessários para reconhecer cada mensagem publicada. Isso controla as garantias de durabilidade que o produtor requer. 0, 1, all (ou -1) todos (Kafka 3.0 e mais recente) 1 (anterior a Kafka 3.0)
Novas tentativas O número de vezes que o cliente reenvia uma mensagem quando o envio encontra um erro. 0,... 0
max.block.ms O número de milissegundos que uma solicitação de envio ou de metadados pode ser bloqueada durante a espera. 0,... 60000 (1 minuto)
max.in.flight.requests.per.connection O número máximo de solicitações não reconhecidas que o cliente envia em uma conexão antes de bloquear outras solicitações. 1,... 5
request.timeout.ms A quantia máxima de tempo que o produtor aguarda por uma resposta a uma solicitação. Se a resposta não for recebida antes do término do tempo limite, a solicitação será tentada novamente ou falhará se o número de tentativas tiver se esgotado. 0,... 30000 (30 segundos)

Há muitas outras definições de configuração disponíveis, mas leia atentamente a documentaçãoApache Kafka antes de experimentá-las.

Particionamento

Com o Kafka, as partições são a unidade de escalabilidade. Portanto, o particionamento é uma maneira eficaz de aumentar sua taxa de transferência, pois permite que os dados do tópico fluam em vários fluxos paralelos.

Quando o produtor publica uma mensagem em um tópico, o produtor pode escolher qual partição usar. Se a ordenação for importante, observe que uma partição é uma sequência ordenada de registros, mas um tópico compreende uma ou mais partições. Se você deseja que um conjunto de mensagens seja entregue em ordem, assegure-se de que todas elas sejam encaminhadas para a mesma partição. A maneira mais direta de fazer isso é fornecer a mesma chave a todas essas mensagens.

O produtor pode especificar explicitamente um número de partição ao publicar uma mensagem. Isso fornece um controle direto, mas torna o código do produtor mais complexo, já que ele assume a responsabilidade de gerenciar a seleção de partição. Para obter mais informações, consulte a chamada de método Producer.partitionsFor. Por exemplo, a chamada é descrita para Kafka versão 2.2.0.

Se o produtor não especificar um número de partição, a seleção de partição será feita por um particionador. O particionador padrão que é construído no produtor Kafka funciona da seguinte forma:

  • Se o registro não tiver uma chave, selecione a partição em um modo round-robin.
  • Se o registro tiver uma chave, selecione a partição calculando um valor do hash para a chave. Isso tem o efeito de selecionar a mesma partição para todas as mensagens com a mesma chave.

Também é possível gravar seu próprio particionador customizado. Um particionador customizado pode escolher qualquer esquema para designar registros às partições. Por exemplo, use apenas um subconjunto das informações na chave ou um identificador específico do aplicativo.

Ordenação de mensagens

O Kafka geralmente grava mensagens na ordem em que elas são enviadas pelo produtor. No entanto, em determinadas situações, as novas tentativas podem fazer com que as mensagens sejam duplicadas ou reordenadas. Se você quiser que uma sequência de mensagens seja enviada em ordem, é importante garantir que todas elas sejam gravadas na mesma partição, pois essa é a única maneira de garantir a ordenação das mensagens.

O produtor também pode tentar enviar mensagens novamente de forma automática. É uma boa ideia habilitar esse recurso de repetição, pois a alternativa é que o código do aplicativo deve executar todas as tentativas por conta própria. A combinação de envio em lote no Kafka e de novas tentativas automáticas pode ter o efeito de duplicação ou de reordenação de mensagens.

Por exemplo, se você publicar uma sequência de três mensagens <M1, M2, M3> em um tópico. Os registros podem caber todos no mesmo lote, portanto, todos são enviados juntos para o líder da partição. O líder, então, os grava na partição e os replica como registros separados. Se ocorrer uma falha, é possível que M1 e M2 sejam adicionados à partição, mas o M3 não é. O produtor não recebe uma confirmação, portanto, tenta novamente enviar <M1, M2, M3>. O novo líder grava M1, M2 e M3 na partição, que agora contém <M1, M2, M1, M2, M3>, onde o M1 duplicado segue o M2 original. Se você restringir o número de solicitações em andamento de cada broker para um só, essa reordenação poderá ser evitada. Você ainda pode descobrir que um único registro é duplicado como <M1, M2, M2, M3>, mas você nunca sai de sequências de pedidos. No Kafka versão 0.11 ou posterior, você também pode usar o recurso de produtor idempotente para evitar a duplicação de M2.

É prática normal do Kafka escrever os aplicativos para lidar com duplicatas ocasionais de mensagens, pois o impacto no desempenho de ter apenas uma única solicitação em andamento é significativo.

Confirmações de mensagem

Ao publicar uma mensagem, você pode escolher o nível de confirmações necessárias usando a configuração do produtor acks. A escolha representa um equilíbrio entre rendimento e confiabilidade. Os três níveis seguintes existem.

acks= 0 (menos confiável)
A mensagem é considerada enviada assim que é gravada na rede. Não há nenhuma confirmação por parte do líder da partição. Como resultado, as mensagens poderão ser perdidas se a liderança da partição mudar. Esse nível de reconhecimento é rápido, mas tem a possibilidade de perda de mensagens em algumas situações.
acks= 1 (o padrão)
A mensagem é confirmada para o produtor assim que o líder da partição escreve com sucesso seu registro na partição. Como a confirmação ocorre antes de o registro chegar às réplicas sincronizadas, a mensagem poderá ser perdida se o líder falhar, mas os seguidores ainda não tiverem a mensagem. Se a liderança da partição mudar, o líder antigo informa o produtor, que pode tratar o erro e tentar enviar novamente a mensagem para o novo líder. Como as mensagens são confirmadas antes de seu recebimento ser confirmado por todas as réplicas, as mensagens que foram confirmadas, mas ainda não foram totalmente replicadas, podem ser perdidas se a liderança da partição mudar.
acks=all (mais confiável)
A mensagem é confirmada para o produtor quando o líder da partição escreve seu registro com sucesso e todas as réplicas sincronizadas fazem o mesmo. A mensagem não será perdida se a liderança da partição mudar, desde que pelo menos uma réplica sincronizada esteja disponível.

Mesmo se você não espera que as mensagens sejam reconhecidas pelo produtor, as mensagens ainda estão disponíveis para serem consumidas somente quando confirmadas, o que significa que a replicação das réplicas sincronizadas está concluída. Em outras palavras, a latência de envio das mensagens do ponto de vista do produtor é menor que a latência de ponta a ponta medida entre o produtor que envia uma mensagem e um consumidor que recebe a mensagem.

Se possível, evite esperar pela confirmação de uma mensagem antes de publicar a próxima mensagem. A espera evita que o produtor envie mensagens em lote juntas e também reduz a taxa em que as mensagens podem ser publicadas abaixo da latência de roundtrip da rede.

Envio em lote, regulagem e compactação

Para fins de eficiência, o produtor coleta lotes de registros para enviar aos servidores. Se você ativar a compactação, o produtor compactará cada lote, o que poderá melhorar o desempenho exigindo que menos dados sejam transferidos pela rede.

Se você tentar publicar mensagens mais rapidamente do que elas podem ser enviadas para um servidor, o produtor as armazenará em buffer automaticamente nas solicitações em lote. O produtor mantém um buffer de registros não enviados para cada partição. Chega um ponto em que nem mesmo a formação de lotes permite que a taxa desejada seja alcançada.

Há outro fator que tem um impacto. Para evitar que produtores ou consumidores individuais sobrecarreguem o cluster, o Event Streams aplica cotas de rendimento. A taxa na qual cada produtor envia dados é calculada e qualquer produtor que tenta exceder sua cota é regulado. A regulagem é aplicada atrasando um pouco o envio de respostas para o produtor. Em geral, isso age apenas como um freio natural.

Para obter mais informações sobre orientação do throughput, consulte Limites e cotas.

Em resumo, quando uma mensagem é publicada, primeiramente seu registro é gravado em um buffer no produtor. Em segundo plano, o produtor envia os registros em lote para o servidor. O servidor, então, responde ao produtor, possivelmente aplicando um atraso de regulagem se o produtor está publicando muito rápido. Se o buffer no produtor ficar cheio, a chamada de envio do produtor será atrasada, mas, em última análise, poderá falhar com uma exceção.

Semântica de entrega

O Kafka oferece a seguinte semântica de entrega de mensagens diferentes:

  • No máximo uma vez: As mensagens podem se perder e não são reenviadas.
  • Pelo menos uma vez: As mensagens nunca são perdidas, mas pode haver duplicatas.
  • Exatamente uma vez: As mensagens nunca são perdidas e não há duplicatas.

A semântica de entrega é determinada pelas configurações a seguir:

  • acks
  • retries
  • enable.idempotence

Por padrão, Kafka usa pelo menos uma semântica.

Para ativar a semântica exatamente uma vez, você deve usar os produtores idempotentes ou transacionais. O produtor idempotente é ativado configurando enable.idempotence para true e garantindo que exatamente uma cópia de cada mensagem seja gravada no Kafka, mesmo que seja tentado novamente. O produtor transacional permite o envio de dados para várias partições, de modo que todas as mensagens sejam entregues com sucesso ou nenhuma delas seja entregue. Ou seja, uma transação é totalmente comprometida ou totalmente descartada. Você também pode incluir offsets em transações para criar aplicativos que leiam, processem e gravem mensagens no Kafka.

Fragmentos de código

Esses fragmentos de código estão em um alto nível para ilustrar os conceitos envolvidos. Para obter exemplos completos, consulte as amostras do Event Streams em GitHub.

Para conectar um consumidor ao Event Streams, deve-se criar credenciais de serviço. Para obter mais informações, consulte Conectando ao Event Streams.

No código do produtor, primeiro é necessário construir o conjunto de propriedades da configuração Todas as conexões com Event Streams são protegidas pelo uso de TLS e autenticação de usuário e senha, portanto, você precisa, no mínimo, dessas propriedades. Substitua BOOTSTRAP_ENDPOINTS, USER e PASSWORD por aqueles de suas próprias credenciais de serviço:

Properties props = new Properties();
 props.put("bootstrap.servers", BOOTSTRAP_ENDPOINTS);
 props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"USER\" password=\"PASSWORD\";");
 props.put("security.protocol", "SASL_SSL");
 props.put("sasl.mechanism", "PLAIN");
 props.put("ssl.protocol", "TLSv1.2");
 props.put("ssl.enabled.protocols", "TLSv1.2");
 props.put("ssl.endpoint.identification.algorithm", "HTTPS");

Para enviar mensagens, você também precisa especificar serializadores para as chaves e os valores, por exemplo:

 props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
 props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

Esses serializadores devem corresponder aos desserializadores usados pelos consumidores

Em seguida, use um KafkaProducer para enviar mensagens, em que cada mensagem é representada por um ProducerRecord. Não se esqueça de fechar o KafkaProducer quando tiver concluído. Esse código apenas envia a mensagem, mas não espera para ver se o envio foi bem-sucedido. A mensagem é enviada para o tópico T1, com chave a sequência key e valor a sequência value.

 Producer<String, String> producer = new KafkaProducer<>(props);
 producer.send(new ProducerRecord<String, String>("T1", "key", "value"));
 producer.close();

O método send() é assíncrono e retorna um Future que pode ser usado para verificar sua conclusão:

 Future<RecordMetadata> f = producer.send(new ProducerRecord<String, String>("T1", "key", "value"));
// Do some other stuff
// Now wait for the result of the send
 RecordMetadata rm = f.get();
 long offset = rm.offset;

Como alternativa, você pode fornecer um retorno de chamada ao enviar a mensagem:

producer.send(new ProducerRecord<String,String>("T1","key","value", new Callback() {
          public void onCompletion(RecordMetadata metadata, Exception exception) {
                     // This is called when the send completes, either successfully or with an exception
          }
});

Para obter mais informações, consulte o Javadoc para o cliente Kafka.{: external}.