Komponent Kafka Client: sgcMQ | eSeGeCe

Kafka Client

TsgcWSPClient_Kafka komunikuje się z Apache Kafka przez własny protokół binarny brokera na porcie 9092. Nie ma przed nim proxy REST ani librdkafka pod spodem, więc producent lub konsument Kafka to pojedynczy, samowystarczalny plik wykonywalny, który możesz debugować od początku do końca.

TsgcWSPClient_Kafka

Kafka to protokół działający na zwykłym TCP, więc nośnikiem jest TsgcTCPClient. Cała reszta jest konfigurowana przez KafkaOptions.

Klasa komponentu

TsgcWSPClient_Kafka

Specyfikacja

Protokół sieciowy Apache Kafka, wsady rekordów v2

Transport

TCP (9092), TLS opcjonalnie

Języki

Delphi, C++ Builder

Transport: zwykły TCP i TLS. sgcMQ łączy się przez zwykły TCP i TLS, a to wszystko, czego wymaga protokół sieciowy Kafka. Pozostałe protokoły sgcMQ definiują transporty WebSocket, na przykład MQTT przez WebSocket lub AMQP przez WebSocket, a te wymagają pakietu sgcWebSockets, który dostarcza klienta WebSocket.

Opublikuj rekord, a potem odpytuj o rekordy

Ustaw GroupId, a komponent dołączy do grupy konsumentów, otrzyma przydział partycji i będzie zatwierdzał offsety za ciebie. Pozostaw je puste, aby czytać partycje bezpośrednio, całkiem bez grupy.

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

Kluczowe właściwości i metody

Składowe, po które sięgasz najczęściej.

Publikowanie

Produce(topic, value, key, partition) zapisuje jeden rekord, ProduceBytes przyjmuje TBytes dla klucza i wartości, a ProduceMessages wysyła cały wsad do jednej partycji i zwraca odpowiedź brokera.

Opcje producenta

KafkaOptions.Producer.Acks wybiera kafkaAcksNone, kafkaAcksLeader lub kafkaAcksAll. Compression włącza gzip, a TimeoutMs ogranicza czas oczekiwania na potwierdzenie brokera.

Konsumowanie

Subscribe([topics]) rejestruje zainteresowanie, a Poll(timeoutMs) pobiera dane. Poll zwraca listę TsgcKafkaMessages, której jesteś właścicielem i którą musisz zwolnić, a dodatkowo zgłasza OnKafkaMessage dla każdego rekordu.

Grupy konsumentów

Ustawienie Consumer.GroupId włącza wykrywanie koordynatora, dołączanie i synchronizację, przydział partycji oraz heartbeaty w tle. OnKafkaRebalance raportuje każdą zmianę przydziału.

Offsety

CommitSync i CommitOffset zatwierdzają jawnie, albo ustaw Consumer.AutoCommit z interwałem. GetEarliestOffset, GetLatestOffset i GetCommittedOffset odczytują pozycje z powrotem.

Bezpośredni odczyt partycji

FetchMessages(topic, partition, offset, maxBytes) odczytuje konkretną partycję od konkretnego offsetu, całkowicie omijając grupę konsumentów.

Strojenie pobierania

Consumer.MinBytes, MaxBytes, MaxPartitionBytes i MaxWaitMs wyważają opóźnienie względem rozmiaru wsadu, a SessionTimeoutMs i RebalanceTimeoutMs rządzą członkostwem w grupie.

Administracja

CreateTopic przyjmuje liczbę partycji i współczynnik replikacji, DeleteTopic usuwa temat, a GetMetadata, ListGroups i DescribeGroups pozwalają zbadać klaster.

Możliwości brokera

GetApiVersions pyta brokera, które wersje API protokołu obsługuje, co jest pewnym sposobem rozgałęzienia kodu według wersji serwera, zamiast zgadywania.

Poznawaj dalej

Pomoc onlinePełna dokumentacja API i przewodnik użytkowania.
All sgcMQ ComponentsPrzejrzyj pełną matrycę funkcji wszystkich siedmiu komponentów.
Pobierz bezpłatną wersję próbnąPublikuj i konsumuj wobec lokalnego brokera Kafka.
CennikLicencje Single, Team i Site z pełnym kodem źródłowym.
Najkorzystniejsza oferta: All-AccessWszystkie produkty eSeGeCe, ze wsparciem Premium w cenie, już od €1,059 rocznie.
Zobacz cennik All-Access

Gotowy, aby zacząć?

Pobierz bezpłatną wersję próbną i opublikuj swój pierwszy rekord Kafka z Delphi lub C++ Builder.