Kafka Client
TsgcWSPClient_Kafka 在 9092 端口上用消息代理自己的二进制协议与 Apache Kafka 通信。它前面没有 REST 代理,下面也没有 librdkafka,因此一个 Kafka 生产者或消费者就是一个自包含的可执行文件,您可以端到端地调试它。
TsgcWSPClient_Kafka 在 9092 端口上用消息代理自己的二进制协议与 Apache Kafka 通信。它前面没有 REST 代理,下面也没有 librdkafka,因此一个 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 接受 TBytes,ProduceMessages 则把整批记录发送到一个分区并返回消息代理的响应。
KafkaOptions.Producer.Acks 可选 kafkaAcksNone、kafkaAcksLeader 或 kafkaAcksAll。Compression 用于开启 gzip,TimeoutMs 则限制消息代理确认的等待时长。
Subscribe([topics]) 注册订阅意向,Poll(timeoutMs) 负责拉取。Poll 返回一个由您拥有并必须释放的 TsgcKafkaMessages 列表,同时为每条记录触发 OnKafkaMessage。
设置 Consumer.GroupId 会启用协调者发现、加入与同步、分区分配以及后台心跳。OnKafkaRebalance 报告每一次分配变化。
CommitSync 和 CommitOffset 用于显式提交,也可以设置带间隔的 Consumer.AutoCommit。GetEarliestOffset、GetLatestOffset 和 GetCommittedOffset 用于读回各个位置。
FetchMessages(topic, partition, offset, maxBytes) 从指定偏移量读取指定分区,完全绕过消费者组。
Consumer.MinBytes、MaxBytes、MaxPartitionBytes 和 MaxWaitMs 用于在延迟和批量大小之间做权衡,SessionTimeoutMs 和 RebalanceTimeoutMs 则决定组成员关系。
CreateTopic 接受分区数和副本因子,DeleteTopic 删除主题,GetMetadata、ListGroups 和 DescribeGroups 则用于检视集群。
GetApiVersions 询问消息代理支持哪些协议 API 版本,这是根据服务器版本分支处理的可靠方式,而不必靠猜。