Kafka クライアントコンポーネント: sgcMQ | eSeGeCe

Kafka クライアント

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 がこれに当たります。それらを利用するには、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.AckskafkaAcksNonekafkaAcksLeaderkafkaAcksAll を選択します。Compression で gzip を有効にし、TimeoutMs がブローカーの ack 待ち時間を制限します。

コンシューム

Subscribe([topics]) で購読を登録し、Poll(timeoutMs) で取得します。Poll は呼び出し側が所有し解放すべき TsgcKafkaMessages のリストを返すとともに、レコードごとに OnKafkaMessage を発生させます。

コンシューマーグループ

Consumer.GroupId を設定すると、コーディネーターの検出、参加と同期、パーティションの割り当て、バックグラウンドのハートビートが有効になります。割り当ての変更は OnKafkaRebalance が通知します。

オフセット

CommitSyncCommitOffset で明示的にコミットするか、Consumer.AutoCommit に間隔を設定します。GetEarliestOffsetGetLatestOffsetGetCommittedOffset で位置を読み出せます。

パーティションの直接読み取り

FetchMessages(topic, partition, offset, maxBytes) は、コンシューマーグループを介さずに、指定したパーティションを指定したオフセットから読み取ります。

フェッチのチューニング

Consumer.MinBytesMaxBytesMaxPartitionBytesMaxWaitMs でレイテンシとバッチサイズのバランスを取り、SessionTimeoutMsRebalanceTimeoutMs がグループのメンバーシップを制御します。

管理操作

CreateTopic はパーティション数とレプリケーションファクターを受け取り、DeleteTopic はトピックを削除します。GetMetadataListGroupsDescribeGroups でクラスターを調べられます。

ブローカーの対応状況

GetApiVersions は、ブローカーが対応しているプロトコル API のバージョンを問い合わせます。サーバーのバージョンを推測せずに分岐するための確実な方法です。

さらに詳しく

オンラインヘルプ完全な API リファレンスと利用ガイドです。
sgcMQ の全コンポーネント7 コンポーネントすべての機能一覧を確認できます。
無料体験版をダウンロードローカルの Kafka ブローカーに対してプロデュースとコンシュームを試せます。
価格Single、Team、Site の各ライセンス。完全なソースコード付きです。
最もお得な選択: All-AccesseSeGeCe の全製品にプレミアムサポートが付いて、年間 €1,059 からご利用いただけます。
All-Access の価格を見る

はじめてみませんか

無料体験版をダウンロードして、Delphi や C++ Builder から最初の Kafka レコードをプロデュースしてみましょう。