Kafka-clientcomponent: sgcMQ | eSeGeCe

Kafka-client

TsgcWSPClient_Kafka praat met Apache Kafka via het eigen binaire protocol van de broker op poort 9092. Er zit geen REST-proxy voor en geen librdkafka onder, dus een Kafka-producer of -consumer is één zelfstandig uitvoerbaar bestand dat je van begin tot eind kunt debuggen.

TsgcWSPClient_Kafka

Kafka is een gewoon TCP-protocol, dus de drager is een TsgcTCPClient. Al het andere stel je in via KafkaOptions.

Componentklasse

TsgcWSPClient_Kafka

Specificatie

Apache Kafka wire protocol, v2-recordbatches

Transport

TCP (9092), TLS optioneel

Talen

Delphi, C++ Builder

Transport: gewoon TCP en TLS. sgcMQ verbindt via gewoon TCP en TLS, en meer vraagt het Kafka wire protocol niet. De andere sgcMQ-protocollen definiëren wel WebSocket-transporten, bijvoorbeeld MQTT over WebSocket of AMQP over WebSocket, en die vereisen een sgcWebSockets-pakket, dat de WebSocket-client levert.

Produceer een record en poll daarna op records

Stel een GroupId in en het component sluit zich aan bij een consumer group, krijgt een partitietoewijzing en commit de offsets voor je. Laat het leeg om partities rechtstreeks te lezen, helemaal zonder group.

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

Belangrijkste eigenschappen & methoden

De leden die je het vaakst gebruikt.

Produceren

Produce(topic, value, key, partition) schrijft één record, ProduceBytes neemt TBytes voor key en value, en ProduceMessages stuurt een hele batch naar één partitie en geeft het antwoord van de broker terug.

Producer-opties

KafkaOptions.Producer.Acks kiest kafkaAcksNone, kafkaAcksLeader of kafkaAcksAll. Compression schakelt gzip in, en TimeoutMs begrenst hoelang de broker op een ack wacht.

Consumeren

Subscribe([topics]) registreert je interesse en Poll(timeoutMs) haalt op. Poll geeft een TsgcKafkaMessages-lijst terug die van jou is en die je moet vrijgeven, en vuurt bovendien OnKafkaMessage per record af.

Consumer groups

Het instellen van Consumer.GroupId schakelt het ontdekken van de coordinator, join en sync, partitietoewijzing en heartbeats op de achtergrond in. OnKafkaRebalance meldt elke wijziging in de toewijzing.

Offsets

CommitSync en CommitOffset committen expliciet, of je stelt Consumer.AutoCommit met een interval in. GetEarliestOffset, GetLatestOffset en GetCommittedOffset lezen posities terug.

Partities rechtstreeks lezen

FetchMessages(topic, partition, offset, maxBytes) leest een specifieke partitie vanaf een specifieke offset en slaat de consumer group volledig over.

Fetch afstemmen

Consumer.MinBytes, MaxBytes, MaxPartitionBytes en MaxWaitMs wegen latentie af tegen batchgrootte, en SessionTimeoutMs en RebalanceTimeoutMs regelen het lidmaatschap van de group.

Beheer

CreateTopic neemt een aantal partities en een replicatiefactor, DeleteTopic verwijdert er een, en GetMetadata, ListGroups en DescribeGroups inspecteren het cluster.

Mogelijkheden van de broker

GetApiVersions vraagt de broker welke protocol-API-versies hij ondersteunt, en dat is de betrouwbare manier om op serverversie te vertakken in plaats van te gokken.

Blijf ontdekken

Online helpVolledige API-referentie en gebruikershandleiding.
Alle sgcMQ-componentenBekijk de volledige functiematrix van alle zeven componenten.
Download de gratis proefversieProduceer en consumeer tegen een lokale Kafka-broker.
PrijzenSingle-, Team- en Site-licenties met volledige broncode.
De beste deal: All-AccessElk eSeGeCe-product, inclusief Premium-ondersteuning, vanaf €1,059 per jaar.
Bekijk de All-Access-prijzen

Klaar om aan de slag te gaan?

Download de gratis proefversie en produceer vanuit Delphi of C++ Builder je eerste Kafka-record.