Kafka コネクタを使用してデータをロードする
StarRocks は、Apache Kafka® コネクタ (StarRocks Connector for Apache Kafka®、以下 Kafka コネクタ) と呼ばれる独自開発のコネクタを提供しています。このコネクタはシンクコネクタとして、Kafka からメッセージを継続的に消費し、それらを StarRocks にロードします。Kafka コネクタは少なくとも一度のセマンティクスを保証します。
Kafka コネクタは Kafka Connect とシームレスに統合でき、これにより StarRocks は Kafka エコシステムとより良く統合されます。リアルタイムデータを StarRocks にロードしたい場合、これは賢明な選択です。Routine Load と比較して、以下のシナリオでは Kafka コネクタの使用が推奨されます。
- Routine Load は CSV、JSON、Avro フォーマットのデータのみをロードできますが、Kafka コネクタは Protobuf など、より多くのフォーマットのデータをロードできます。Kafka Connect のコンバータを使用してデータを JSON や CSV フォーマットに変換できる限り、データは Kafka コネクタを介して StarRocks にロードできます。
- Debezium フォーマットの CDC データなど、データ変換をカスタマイズする。
- 複数の Kafka トピックからデータをロードする。
- Confluent Cloud からデータをロードする。
- ロードバッチサイズ、並行性、その他のパラメータを細かく制御して、ロード速度とリソース利用のバランスを取る必要がある。
準備
バージョン要件
| コネクタ | Kafka | StarRocks | Java |
|---|---|---|---|
| 1.0.4 | 3.4 | 2.5 and later | 8 |
| 1.0.3 | 3.4 | 2.5 and later | 8 |
Kafka 環境のセットアップ
自己管理型の Apache Kafka クラスターと Confluent Cloud の両方がサポートされています。
- 自己管理型の Apache Kafka クラスターの場合、Apache Kafka クイックスタート を参照して、Kafka クラスターを迅速にデプロイできます。Kafka Connect はすでに Kafka に統合されています。
- Confluent Cloud の場合、Confluent アカウントを持ち、クラスターを作成していることを確認してください。
Kafka コネクタのダウンロード
Kafka コネクタを Kafka Connect に提出します。
-
自己管理型 Kafka クラスター:
starrocks-kafka-connector-xxx.tar.gz をダウンロードして解凍します。
-
Confluent Cloud:
現在、Kafka コネクタは Confluent Hub にアップロードされていません。starrocks-kafka-connector-xxx.tar.gz をダウンロードして解凍し、ZIP ファイルにパッケージ化して Confluent Cloud にアップロードする必要があります。
ネットワーク設定
Kafka が配置されているマシンが StarRocks クラスターの FE ノードに http_port (デフォルト: 8030) と query_port (デフォルト: 9030) を介してアクセスでき、BE ノードに be_http_port (デフォルト: 8040) を介してアクセスできることを確認してください。
使用法
このセクションでは、自己管理型 Kafka クラスターを例にとり、Kafka コネクタと Kafka Connect の設定方法を説明し、その後 Kafka Connect を 実行して StarRocks にデータをロードする方法を説明します。
データセットの準備
Kafka クラスターのトピック test に JSON フォーマットのデータが存在すると仮定します。
{"id":1,"city":"New York"}
{"id":2,"city":"Los Angeles"}
{"id":3,"city":"Chicago"}
テーブルの作成
StarRocks クラスターのデータベース example_db に JSON フォーマットのデータのキーに従ってテーブル test_tbl を作成します。
CREATE DATABASE example_db;
USE example_db;
CREATE TABLE test_tbl (id INT, city STRING);