Componente Kafka Client: sgcMQ | eSeGeCe

Kafka Client

TsgcWSPClient_Kafka dialoga con Apache Kafka usando il protocollo binario del broker sulla porta 9092. Davanti non c'è nessun proxy REST e sotto non c'è nessuna librdkafka, quindi un producer o un consumer Kafka è un unico eseguibile autosufficiente che puoi seguire nel debug dall'inizio alla fine.

TsgcWSPClient_Kafka

Kafka è un protocollo su TCP semplice, quindi il vettore è un TsgcTCPClient. Tutto il resto si configura attraverso KafkaOptions.

Classe del componente

TsgcWSPClient_Kafka

Specifica

Protocollo binario Apache Kafka, record batch v2

Trasporto

TCP (9092), TLS facoltativo

Linguaggi

Delphi, C++ Builder

Trasporto: TCP semplice e TLS. sgcMQ si connette tramite TCP semplice e TLS, che è tutto ciò che il protocollo binario di Kafka richiede. Gli altri protocolli di sgcMQ definiscono anche trasporti WebSocket, per esempio MQTT su WebSocket o AMQP su WebSocket, e quelli richiedono un package sgcWebSockets, che fornisce il client WebSocket.

Produci un record, poi interroga i record

Imposta un GroupId e il componente entra in un consumer group, riceve un'assegnazione di partizioni e conferma gli offset al posto tuo. Lascialo vuoto per leggere le partizioni direttamente, senza nessun gruppo.

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;
}

Proprietà e metodi principali

I membri che userai più spesso.

Produzione

Produce(topic, value, key, partition) scrive un record, ProduceBytes accetta TBytes per chiave e valore, e ProduceMessages invia un intero batch a una sola partizione e restituisce la risposta del broker.

Opzioni del producer

KafkaOptions.Producer.Acks sceglie kafkaAcksNone, kafkaAcksLeader o kafkaAcksAll. Compression attiva gzip, e TimeoutMs limita l'attesa della conferma del broker.

Consumo

Subscribe([topics]) registra l'interesse e Poll(timeoutMs) preleva. Poll restituisce una lista TsgcKafkaMessages di cui sei proprietario e che devi liberare, e solleva anche OnKafkaMessage per ogni record.

Consumer group

Impostare Consumer.GroupId attiva la scoperta del coordinatore, il join e la sincronizzazione, l'assegnazione delle partizioni e gli heartbeat in background. OnKafkaRebalance segnala ogni cambio di assegnazione.

Offset

CommitSync e CommitOffset confermano esplicitamente, oppure imposta Consumer.AutoCommit con un intervallo. GetEarliestOffset, GetLatestOffset e GetCommittedOffset rileggono le posizioni.

Lettura diretta delle partizioni

FetchMessages(topic, partition, offset, maxBytes) legge una partizione specifica a partire da un offset specifico, aggirando del tutto il consumer group.

Regolazione del fetch

Consumer.MinBytes, MaxBytes, MaxPartitionBytes e MaxWaitMs bilanciano la latenza contro la dimensione del batch, mentre SessionTimeoutMs e RebalanceTimeoutMs regolano l'appartenenza al gruppo.

Amministrazione

CreateTopic accetta un numero di partizioni e un fattore di replica, DeleteTopic ne rimuove uno, e GetMetadata, ListGroups e DescribeGroups ispezionano il cluster.

Capability del broker

GetApiVersions chiede al broker quali versioni delle API del protocollo supporta, che è il modo affidabile per differenziare il comportamento in base alla versione del server invece di tirare a indovinare.

Continua a esplorare

Guida onlineRiferimento API completo e guida all'utilizzo.
Tutti i componenti sgcMQEsplora la matrice completa delle funzionalità di tutti e sette i componenti.
Scarica la versione di prova gratuitaProduci e consuma contro un broker Kafka locale.
PrezziLicenze Single, Team e Site con codice sorgente completo.
La scelta più conveniente: All-AccessTutti i prodotti eSeGeCe, con Supporto Premium incluso, a partire da €1,059/anno.
Vedi i prezzi All-Access

Pronto per iniziare?

Scarica la versione di prova gratuita e produci il tuo primo record Kafka da Delphi o C++ Builder.