Componente cliente Kafka: sgcMQ | eSeGeCe

Cliente Kafka

O TsgcWSPClient_Kafka conversa com o Apache Kafka usando o protocolo binário do próprio broker na porta 9092. Não há um proxy REST na frente nem a librdkafka por baixo, então um produtor ou consumidor Kafka é um único executável autossuficiente que você depura de ponta a ponta.

TsgcWSPClient_Kafka

O Kafka é um protocolo sobre TCP puro, então o transporte é um TsgcTCPClient. Todo o resto é configurado por meio de KafkaOptions.

Classe do componente

TsgcWSPClient_Kafka

Especificação

Protocolo binário do Apache Kafka, record batches v2

Transporte

TCP (9092), TLS opcional

Linguagens

Delphi, C++ Builder

Transporte: TCP puro e TLS. O sgcMQ se conecta por TCP puro e TLS, que é tudo o que o protocolo binário do Kafka pede. Os outros protocolos do sgcMQ definem, sim, transportes WebSocket, por exemplo MQTT sobre WebSocket ou AMQP sobre WebSocket, e esses exigem um pacote sgcWebSockets, que fornece o cliente WebSocket.

Produza um registro e depois faça poll dos registros

Defina um GroupId e o componente entra em um grupo de consumidores, recebe uma atribuição de partições e faz o commit dos offsets por você. Deixe vazio para ler as partições diretamente, sem grupo nenhum.

uses
  sgcTCP_Client_WS, sgcWebSocket_Protocols, sgcKafka_Classes;

var
  TCPClient: TsgcTCPClient;
  Kafka: TsgcWSPClient_Kafka;
begin
  Kafka := TsgcWSPClient_Kafka.Create(nil);
  Kafka.KafkaOptions.ClientId := 'my-delphi-app';
  Kafka.KafkaOptions.Producer.Acks        := kafkaAcksLeader;
  Kafka.KafkaOptions.Producer.Compression := kafkaCompressionGzip;
  Kafka.KafkaOptions.Consumer.GroupId     := 'my-group';
  Kafka.KafkaOptions.Consumer.OffsetReset := kafkaOffsetEarliest;
  Kafka.OnKafkaMessage := KafkaMessage;
  Kafka.OnKafkaProduce := KafkaProduce;

  TCPClient := TsgcTCPClient.Create(nil);
  Kafka.Client := TCPClient;
  TCPClient.Host := '127.0.0.1';
  TCPClient.Port := 9092;
  TCPClient.Active := True;

  // produce a record to a topic
  Kafka.Produce('my-topic', 'Hello Kafka', 'key-1');

  // consume: subscribe once, then Poll repeatedly (e.g. from a timer)
  Kafka.Subscribe(['my-topic']);
end;

procedure TForm1.KafkaMessage(Sender: TObject;
  const Message: TsgcKafkaMessage);
begin
  Memo1.Lines.Add(Message.GetKeyString + ' = ' + Message.GetValueString);
end;

// fetch records, commit when a batch is processed
var
  Messages: TsgcKafkaMessages;
begin
  Messages := Kafka.Poll(1000);
  try
    if Messages.Count > 0 then
      Kafka.CommitSync;
  finally
    Messages.Free;
  end;
end;
// include: sgcTCP_Client_WS.hpp, sgcWebSocket_Protocols.hpp,
// sgcKafka_Classes.hpp

TsgcWSPClient_Kafka *Kafka = new TsgcWSPClient_Kafka(this);
Kafka->KafkaOptions->ClientId = "my-cbuilder-app";
Kafka->KafkaOptions->Producer->Acks = kafkaAcksLeader;
Kafka->KafkaOptions->Consumer->GroupId = "my-group";
Kafka->KafkaOptions->Consumer->OffsetReset = kafkaOffsetEarliest;
Kafka->OnKafkaMessage = KafkaMessage;

TsgcTCPClient *TCPClient = new TsgcTCPClient(this);
Kafka->Client = TCPClient;
TCPClient->Host = "127.0.0.1";
TCPClient->Port = 9092;
TCPClient->Active = true;

Kafka->Produce("my-topic", "Hello Kafka", "key-1");
Kafka->Subscribe(ARRAYOFCONST(("my-topic")));

TsgcKafkaMessages *Messages = Kafka->Poll(1000);
try {
  if (Messages->Count > 0)
    Kafka->CommitSync();
}
__finally {
  delete Messages;
}

Principais propriedades e métodos

Os membros que você usa com mais frequência.

Produção

Produce(topic, value, key, partition) grava um registro, ProduceBytes recebe TBytes para a chave e o valor, e ProduceMessages envia um lote inteiro para uma partição e retorna a resposta do broker.

Opções do produtor

KafkaOptions.Producer.Acks escolhe kafkaAcksNone, kafkaAcksLeader ou kafkaAcksAll. Compression liga o gzip, e TimeoutMs limita a espera pelo ack do broker.

Consumo

Subscribe([topics]) registra o interesse e Poll(timeoutMs) busca os dados. O Poll retorna uma lista TsgcKafkaMessages que pertence a você e precisa ser liberada, e também dispara OnKafkaMessage a cada registro.

Grupos de consumidores

Definir Consumer.GroupId ativa a descoberta do coordenador, o join e o sync, a atribuição de partições e os heartbeats em segundo plano. OnKafkaRebalance informa cada mudança de atribuição.

Offsets

CommitSync e CommitOffset fazem o commit de forma explícita, ou defina Consumer.AutoCommit com um intervalo. GetEarliestOffset, GetLatestOffset e GetCommittedOffset leem as posições de volta.

Leitura direta de partições

FetchMessages(topic, partition, offset, maxBytes) lê uma partição específica a partir de um offset específico, ignorando por completo o grupo de consumidores.

Ajuste do fetch

Consumer.MinBytes, MaxBytes, MaxPartitionBytes e MaxWaitMs equilibram latência e tamanho do lote, e SessionTimeoutMs e RebalanceTimeoutMs governam a participação no grupo.

Administração

CreateTopic recebe a quantidade de partições e o fator de replicação, DeleteTopic remove um tópico, e GetMetadata, ListGroups e DescribeGroups inspecionam o cluster.

Capacidades do broker

GetApiVersions pergunta ao broker quais versões da API do protocolo ele suporta, que é a forma confiável de ramificar por versão do servidor em vez de adivinhar.

Continue explorando

Ajuda onlineReferência completa da API e guia de uso.
Todos os componentes sgcMQNavegue pela matriz completa de recursos dos sete componentes.
Baixar avaliação gratuitaProduza e consuma contra um broker Kafka local.
PreçosLicenças Single, Team e Site com código-fonte completo.
Melhor custo-benefício: All-AccessTodos os produtos da eSeGeCe, com Suporte Premium incluído, a partir de €1,059/ano.
Ver preços do All-Access

Pronto para começar?

Baixe a versão de avaliação gratuita e produza seu primeiro registro Kafka a partir do Delphi ou do C++ Builder.