Componente Kafka Client: sgcMQ | eSeGeCe

Kafka Client

TsgcWSPClient_Kafka habla con Apache Kafka usando el propio protocolo binario del broker en el puerto 9092. No hay un proxy REST por delante ni librdkafka por debajo, así que un productor o un consumidor de Kafka es un único ejecutable autocontenido que puedes depurar de principio a fin.

TsgcWSPClient_Kafka

Kafka es un protocolo de TCP plano, así que el transporte es un TsgcTCPClient. Todo lo demás se configura desde KafkaOptions.

Clase del componente

TsgcWSPClient_Kafka

Especificación

Protocolo binario de Apache Kafka, lotes de registros v2

Transporte

TCP (9092), TLS opcional

Lenguajes

Delphi, C++ Builder

Transporte: TCP plano y TLS. sgcMQ se conecta por TCP plano y TLS, que es todo lo que pide el protocolo binario de Kafka. Los demás protocolos de sgcMQ sí definen transportes WebSocket, por ejemplo MQTT sobre WebSocket o AMQP sobre WebSocket, y esos requieren un paquete sgcWebSockets, que es el que aporta el cliente WebSocket.

Produce un registro y luego haz poll de registros

Si defines un GroupId, el componente se une a un grupo de consumidores, recibe una asignación de particiones y hace los commits de offsets por ti. Déjalo vacío para leer particiones directamente, sin grupo alguno.

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

Propiedades y métodos principales

Los miembros que vas a usar con más frecuencia.

Producción

Produce(topic, value, key, partition) escribe un registro, ProduceBytes recibe TBytes para la clave y el valor, y ProduceMessages envía un lote completo a una partición y devuelve la respuesta del broker.

Opciones del productor

KafkaOptions.Producer.Acks elige entre kafkaAcksNone, kafkaAcksLeader y kafkaAcksAll. Compression activa gzip, y TimeoutMs limita la espera del ack del broker.

Consumo

Subscribe([topics]) registra el interés y Poll(timeoutMs) recoge los datos. Poll devuelve una lista TsgcKafkaMessages que te pertenece y debes liberar, y además dispara OnKafkaMessage por cada registro.

Grupos de consumidores

Definir Consumer.GroupId activa el descubrimiento del coordinador, el join y el sync, la asignación de particiones y los heartbeats en segundo plano. OnKafkaRebalance informa de cada cambio de asignación.

Offsets

CommitSync y CommitOffset confirman de forma explícita, o activa Consumer.AutoCommit con un intervalo. GetEarliestOffset, GetLatestOffset y GetCommittedOffset vuelven a leer las posiciones.

Lecturas directas de particiones

FetchMessages(topic, partition, offset, maxBytes) lee una partición concreta desde un offset concreto, saltándose por completo el grupo de consumidores.

Ajuste del fetch

Consumer.MinBytes, MaxBytes, MaxPartitionBytes y MaxWaitMs equilibran latencia y tamaño de lote, y SessionTimeoutMs y RebalanceTimeoutMs rigen la pertenencia al grupo.

Administración

CreateTopic recibe un número de particiones y un factor de replicación, DeleteTopic elimina uno, y GetMetadata, ListGroups y DescribeGroups inspeccionan el clúster.

Capacidades del broker

GetApiVersions pregunta al broker qué versiones de la API del protocolo admite, que es la forma fiable de ramificar según la versión del servidor en lugar de adivinarla.

Sigue explorando

Ayuda en líneaReferencia completa de la API y guía de uso.
Todos los componentes de sgcMQConsulta la matriz completa de características de los siete componentes.
Descargar prueba gratuitaProduce y consume contra un broker Kafka local.
PreciosLicencias Single, Team y Site con código fuente completo.
La mejor opción: All-AccessTodos los productos de eSeGeCe, con Premium Support incluido, desde €1,059 al año.
Ver precios de All-Access

¿Listo para empezar?

Descarga la prueba gratuita y produce tu primer registro de Kafka desde Delphi o C++ Builder.