Kafka-Client-Komponente: sgcMQ | eSeGeCe

Kafka-Client

TsgcWSPClient_Kafka spricht mit Apache Kafka über das brokereigene Binärprotokoll auf Port 9092. Davor steht kein REST-Proxy und darunter liegt kein librdkafka, ein Kafka-Producer oder -Consumer ist also eine einzige, in sich geschlossene ausführbare Datei, die du von Anfang bis Ende debuggen kannst.

TsgcWSPClient_Kafka

Kafka ist ein reines TCP-Protokoll, der Träger ist also ein TsgcTCPClient. Alles Weitere konfigurierst du über KafkaOptions.

Komponentenklasse

TsgcWSPClient_Kafka

Spezifikation

Apache Kafka Wire-Protokoll, Record Batches v2

Transport

TCP (9092), TLS optional

Sprachen

Delphi, C++ Builder

Transport: reines TCP und TLS. sgcMQ verbindet sich über reines TCP und TLS, und mehr verlangt das Kafka-Wire-Protokoll auch nicht. Die anderen sgcMQ-Protokolle definieren sehr wohl WebSocket-Transporte, zum Beispiel MQTT über WebSocket oder AMQP über WebSocket, und dafür benötigst du ein sgcWebSockets Paket, das den WebSocket-Client bereitstellt.

Einen Record erzeugen, dann Records abholen

Setzt du eine GroupId, tritt die Komponente einer Consumer Group bei, erhält eine Partitionszuteilung und committet die Offsets für dich. Lässt du sie leer, liest sie Partitionen direkt, ganz ohne Gruppe.

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

Wichtige Eigenschaften & Methoden

Die Member, zu denen du am häufigsten greifst.

Erzeugen

Produce(topic, value, key, partition) schreibt einen einzelnen Record, ProduceBytes nimmt TBytes für Key und Value entgegen, und ProduceMessages sendet einen ganzen Batch an eine Partition und gibt die Broker-Antwort zurück.

Producer-Optionen

KafkaOptions.Producer.Acks wählt kafkaAcksNone, kafkaAcksLeader oder kafkaAcksAll. Compression schaltet gzip ein, und TimeoutMs begrenzt, wie lange der Broker auf das Ack wartet.

Konsumieren

Subscribe([topics]) meldet das Interesse an und Poll(timeoutMs) holt ab. Poll gibt eine TsgcKafkaMessages-Liste zurück, die dir gehört und die du freigeben musst, und löst zusätzlich pro Record OnKafkaMessage aus.

Consumer Groups

Setzt du Consumer.GroupId, aktiviert das die Coordinator-Suche, Join und Sync, die Partitionszuteilung und Heartbeats im Hintergrund. OnKafkaRebalance meldet jede Änderung der Zuteilung.

Offsets

CommitSync und CommitOffset committen explizit, alternativ setzt du Consumer.AutoCommit mit einem Intervall. GetEarliestOffset, GetLatestOffset und GetCommittedOffset lesen Positionen wieder aus.

Direkte Partitionslesevorgänge

FetchMessages(topic, partition, offset, maxBytes) liest eine bestimmte Partition ab einem bestimmten Offset und umgeht die Consumer Group vollständig.

Fetch-Feinabstimmung

Consumer.MinBytes, MaxBytes, MaxPartitionBytes und MaxWaitMs wägen Latenz gegen Batch-Größe ab, und SessionTimeoutMs und RebalanceTimeoutMs steuern die Gruppenmitgliedschaft.

Verwaltung

CreateTopic nimmt eine Partitionsanzahl und einen Replikationsfaktor entgegen, DeleteTopic entfernt ein Topic, und GetMetadata, ListGroups und DescribeGroups inspizieren den Cluster.

Broker-Fähigkeiten

GetApiVersions fragt den Broker, welche Protokoll-API-Versionen er unterstützt. Das ist der zuverlässige Weg, um nach Serverversion zu verzweigen, statt zu raten.

Weiter entdecken

Online-HilfeVollständige API-Referenz und Anwendungsleitfaden.
Alle sgcMQ-KomponentenDurchstöbere die vollständige Funktionsmatrix aller sieben Komponenten.
Kostenlose Testversion herunterladenErzeuge und konsumiere gegen einen lokalen Kafka-Broker.
PreiseSingle-, Team- und Site-Lizenzen mit vollständigem Quellcode.
Bestes Preis-Leistungs-Verhältnis: All-AccessAlle eSeGeCe-Produkte, inklusive Premium-Support, ab €1,059 pro Jahr.
All-Access-Preise ansehen

Bereit loszulegen?

Lade die kostenlose Testversion herunter und erzeuge deinen ersten Kafka-Record aus Delphi oder C++ Builder.