Kafka İstemci Bileşeni: sgcMQ | eSeGeCe

Kafka İstemcisi

TsgcWSPClient_Kafka, Apache Kafka ile 9092 numaralı bağlantı noktasında broker'ın kendi ikili protokolü üzerinden konuşur. Önünde bir REST proxy'si, altında bir librdkafka yoktur; böylece bir Kafka üreticisi ya da tüketicisi, uçtan uca hata ayıklayabileceğiniz tek ve kendi kendine yeten bir yürütülebilir dosyadır.

TsgcWSPClient_Kafka

Kafka düz bir TCP protokolüdür, bu nedenle taşıyıcı bir TsgcTCPClient'tır. Geri kalan her şey KafkaOptions üzerinden yapılandırılır.

Bileşen sınıfı

TsgcWSPClient_Kafka

Spesifikasyon

Apache Kafka hat protokolü, v2 kayıt yığınları

Taşıma

TCP (9092), isteğe bağlı TLS

Diller

Delphi, C++ Builder

Taşıma: düz TCP ve TLS. sgcMQ, düz TCP ve TLS üzerinden bağlanır; Kafka hat protokolünün istediği de zaten budur. Diğer sgcMQ protokolleri WebSocket taşımaları tanımlar, örneğin WebSocket üzerinden MQTT ya da WebSocket üzerinden AMQP; bunlar ise WebSocket istemcisini sağlayan bir sgcWebSockets paketi gerektirir.

Bir kayıt üretin, ardından kayıtları yoklayın

Bir GroupId ayarlarsanız bileşen bir tüketici grubuna katılır, bir bölüm ataması alır ve offset'leri sizin için kesinleştirir. Hiçbir grup olmadan bölümleri doğrudan okumak için bunu boş bırakın.

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

Temel özellikler & yöntemler

En sık başvurduğunuz üyeler.

Üretme

Produce(topic, value, key, partition) tek bir kayıt yazar, ProduceBytes anahtar ve değer için TBytes alır, ProduceMessages ise bütün bir yığını tek bir bölüme gönderir ve broker yanıtını döndürür.

Üretici seçenekleri

KafkaOptions.Producer.Acks; kafkaAcksNone, kafkaAcksLeader ya da kafkaAcksAll arasından seçim yapar. Compression gzip'i devreye alır, TimeoutMs ise broker'ın onay beklemesini sınırlar.

Tüketme

Subscribe([topics]) ilgi kaydı oluşturur, Poll(timeoutMs) ise getirir. Poll, sahipliği sizde olan ve serbest bırakmanız gereken bir TsgcKafkaMessages listesi döndürür, ayrıca her kayıt için OnKafkaMessage olayını tetikler.

Tüketici grupları

Consumer.GroupId ayarlamak; koordinatör keşfini, katılma ve eşitlemeyi, bölüm atamasını ve arka plan heartbeat'lerini açar. OnKafkaRebalance her atama değişikliğini raporlar.

Offset'ler

CommitSync ve CommitOffset açıkça kesinleştirir, ya da bir aralıkla birlikte Consumer.AutoCommit ayarlarsınız. GetEarliestOffset, GetLatestOffset ve GetCommittedOffset konumları geri okur.

Doğrudan bölüm okumaları

FetchMessages(topic, partition, offset, maxBytes), tüketici grubunu tamamen atlayarak belirli bir bölümü belirli bir offset'ten okur.

Getirme ayarı

Consumer.MinBytes, MaxBytes, MaxPartitionBytes ve MaxWaitMs, gecikmeyi yığın boyutuna karşı dengeler; SessionTimeoutMs ve RebalanceTimeoutMs ise grup üyeliğini yönetir.

Yönetim

CreateTopic bir bölüm sayısı ve çoğaltma faktörü alır, DeleteTopic bir konuyu kaldırır; GetMetadata, ListGroups ve DescribeGroups ise kümeyi inceler.

Broker yetenekleri

GetApiVersions, broker'a hangi protokol API sürümlerini desteklediğini sorar; sunucu sürümüne göre dallanmanın tahmin yürütmek yerine güvenilir yolu budur.

Keşfetmeye devam edin

Çevrimiçi yardımTam API referansı ve kullanım kılavuzu.
Tüm sgcMQ BileşenleriYedi bileşenin tam özellik matrisine göz atın.
Ücretsiz Deneme Sürümünü İndirinYerel bir Kafka broker'ına karşı üretin ve tüketin.
FiyatlandırmaTam kaynak kodlu Single, Team ve Site lisansları.
En avantajlı seçenek: All-AccessTüm eSeGeCe ürünleri, Premium Destek dahil, yılda €1,059'dan itibaren.
All-Access fiyatlarına bakın

Başlamaya Hazır mısınız?

Ücretsiz deneme sürümünü indirin ve Delphi ya da C++ Builder'dan ilk Kafka kaydınızı üretin.