Composant client Kafka : sgcMQ | eSeGeCe

Client Kafka

TsgcWSPClient_Kafka dialogue avec Apache Kafka via le protocole binaire propre au broker, sur le port 9092. Il n'y a aucun proxy REST devant lui ni librdkafka en dessous, un producteur ou un consommateur Kafka est donc un seul exécutable autonome que tu peux déboguer de bout en bout.

TsgcWSPClient_Kafka

Kafka est un protocole en TCP simple, le support de transport est donc un TsgcTCPClient. Tout le reste se configure via KafkaOptions.

Classe du composant

TsgcWSPClient_Kafka

Spécification

Protocole binaire Apache Kafka, lots d'enregistrements v2

Transport

TCP (9092), TLS facultatif

Langages

Delphi, C++ Builder

Transport : TCP simple et TLS. sgcMQ se connecte en TCP simple et en TLS, ce qui est tout ce que demande le protocole binaire de Kafka. Les autres protocoles de sgcMQ définissent bien des transports WebSocket, par exemple MQTT sur WebSocket ou AMQP sur WebSocket, et ceux-là nécessitent un package sgcWebSockets, qui fournit le client WebSocket.

Produire un enregistrement, puis interroger

Renseigne un GroupId et le composant rejoint un groupe de consommateurs, reçoit une affectation de partitions et valide les offsets à ta place. Laisse-le vide pour lire les partitions directement, sans aucun groupe.

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

Propriétés et méthodes clés

Les membres que tu utilises le plus souvent.

Production

Produce(topic, value, key, partition) écrit un enregistrement, ProduceBytes accepte des TBytes pour la clé et la valeur, et ProduceMessages envoie un lot entier vers une partition et renvoie la réponse du broker.

Options du producteur

KafkaOptions.Producer.Acks choisit kafkaAcksNone, kafkaAcksLeader ou kafkaAcksAll. Compression active gzip, et TimeoutMs borne l'attente d'acquittement du broker.

Consommation

Subscribe([topics]) déclare l'intérêt et Poll(timeoutMs) récupère. Poll renvoie une liste TsgcKafkaMessages qui t'appartient et que tu dois libérer, et déclenche aussi OnKafkaMessage par enregistrement.

Groupes de consommateurs

Renseigner Consumer.GroupId active la découverte du coordinateur, l'adhésion et la synchronisation, l'affectation des partitions et les heartbeats en arrière-plan. OnKafkaRebalance signale chaque changement d'affectation.

Offsets

CommitSync et CommitOffset valident explicitement, ou active Consumer.AutoCommit avec un intervalle. GetEarliestOffset, GetLatestOffset et GetCommittedOffset relisent les positions.

Lectures directes de partition

FetchMessages(topic, partition, offset, maxBytes) lit une partition précise depuis un offset précis, en contournant entièrement le groupe de consommateurs.

Réglage de la récupération

Consumer.MinBytes, MaxBytes, MaxPartitionBytes et MaxWaitMs arbitrent entre latence et taille de lot, et SessionTimeoutMs et RebalanceTimeoutMs régissent l'appartenance au groupe.

Administration

CreateTopic prend un nombre de partitions et un facteur de réplication, DeleteTopic en supprime un, et GetMetadata, ListGroups et DescribeGroups inspectent le cluster.

Capacités du broker

GetApiVersions demande au broker quelles versions d'API du protocole il prend en charge, c'est la façon fiable de brancher selon la version du serveur au lieu de deviner.

Pour aller plus loin

Aide en ligneRéférence complète de l'API et guide d'utilisation.
Tous les composants sgcMQParcours la matrice complète des fonctionnalités des sept composants.
Télécharger l'essai gratuitProduis et consomme contre un broker Kafka local.
TarifsLicences Single, Team et Site avec le code source complet.
Meilleur rapport qualité-prix : All-AccessTous les produits eSeGeCe, Support Premium inclus, à partir de €1,059/an.
Voir les tarifs All-Access

Prêt à te lancer ?

Télécharge l'essai gratuit et produis ton premier enregistrement Kafka depuis Delphi ou C++ Builder.