Kafka クライアント
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 がこれに当たります。それらを利用するには、WebSocket クライアントを提供する sgcWebSockets のパッケージが必要です。
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) は 1 件のレコードを書き込みます。ProduceBytes はキーと値に TBytes を受け取り、ProduceMessages はバッチ全体を 1 つのパーティションへ送信してブローカーの応答を返します。
KafkaOptions.Producer.Acks で kafkaAcksNone、kafkaAcksLeader、kafkaAcksAll を選択します。Compression で gzip を有効にし、TimeoutMs がブローカーの ack 待ち時間を制限します。
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 のバージョンを問い合わせます。サーバーのバージョンを推測せずに分岐するための確実な方法です。