Kafka Client 组件:sgcMQ | eSeGeCe

Kafka Client

TsgcWSPClient_Kafka 在 9092 端口上用消息代理自己的二进制协议与 Apache Kafka 通信。它前面没有 REST 代理,下面也没有 librdkafka,因此一个 Kafka 生产者或消费者就是一个自包含的可执行文件,您可以端到端地调试它。

TsgcWSPClient_Kafka

Kafka 是纯 TCP 协议,因此载体就是一个 TsgcTCPClient。其余一切都通过 KafkaOptions 配置。

组件类

TsgcWSPClient_Kafka

规范

Apache Kafka 线路协议,v2 记录批

传输

TCP (9092),TLS 可选

语言

Delphi、C++ Builder

传输:纯 TCP 和 TLS。sgcMQ 通过纯 TCP 和 TLS 连接,而 Kafka 线路协议所需要的也正是如此。sgcMQ 的其他协议确实定义了 WebSocket 传输,例如 MQTT over WebSocket 或 AMQP over WebSocket,那些场景需要 sgcWebSockets 包,它提供 WebSocket 客户端。

先生产一条记录,再轮询获取记录

设置 GroupId,组件就会加入一个消费者组、接收分区分配并代您提交偏移量。留空则完全不使用消费者组,直接读取分区。

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

关键属性与方法

您最常使用的成员。

生产

Produce(topic, value, key, partition) 写入一条记录,ProduceBytes 的 key 和 value 接受 TBytesProduceMessages 则把整批记录发送到一个分区并返回消息代理的响应。

生产者选项

KafkaOptions.Producer.Acks 可选 kafkaAcksNonekafkaAcksLeaderkafkaAcksAllCompression 用于开启 gzip,TimeoutMs 则限制消息代理确认的等待时长。

消费

Subscribe([topics]) 注册订阅意向,Poll(timeoutMs) 负责拉取。Poll 返回一个由您拥有并必须释放的 TsgcKafkaMessages 列表,同时为每条记录触发 OnKafkaMessage

消费者组

设置 Consumer.GroupId 会启用协调者发现、加入与同步、分区分配以及后台心跳。OnKafkaRebalance 报告每一次分配变化。

偏移量

CommitSyncCommitOffset 用于显式提交,也可以设置带间隔的 Consumer.AutoCommitGetEarliestOffsetGetLatestOffsetGetCommittedOffset 用于读回各个位置。

直接读取分区

FetchMessages(topic, partition, offset, maxBytes) 从指定偏移量读取指定分区,完全绕过消费者组。

拉取调优

Consumer.MinBytesMaxBytesMaxPartitionBytesMaxWaitMs 用于在延迟和批量大小之间做权衡,SessionTimeoutMsRebalanceTimeoutMs 则决定组成员关系。

管理

CreateTopic 接受分区数和副本因子,DeleteTopic 删除主题,GetMetadataListGroupsDescribeGroups 则用于检视集群。

消息代理能力

GetApiVersions 询问消息代理支持哪些协议 API 版本,这是根据服务器版本分支处理的可靠方式,而不必靠猜。

继续探索

在线帮助完整的 API 参考和使用指南。
所有 sgcMQ 组件浏览全部七个组件的完整功能矩阵。
下载免费试用版对接本地 Kafka 消息代理进行生产和消费。
价格Single、Team 和 Site 授权,均含完整源代码。
超值之选:All-AccesseSeGeCe 全部产品,含高级支持,每年 €1,059 起。
查看 All-Access 价格

准备好开始了吗?

下载免费试用版,从 Delphi 或 C++ Builder 生产您的第一条 Kafka 记录。