Kafka コネクタを使用してデータをロードする
StarRocks は、Apache Kafka® コネクタ (StarRocks Connector for Apache 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);
Kafka コネクタと Kafka Connect の設定と実行、データのロード
スタンドアロンモードで Kafka Connect を実行
-
Kafka コネクタを設定します。Kafka インストールディレクトリの config ディレクトリに、Kafka コネクタ用の設定ファイル connect-StarRocks-sink.properties を作成し、以下のパラメータを設定します。詳細なパラメータと説明については、Parameters を参照してください。
注記Kafka コネクタはシンクコネクタです。
name=starrocks-kafka-connector
connector.class=com.starrocks.connector.kafka.StarRocksSinkConnector
topics=test
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=false
# StarRocks クラスター内の FE の HTTP URL。デフォルトポートは 8030 です。
starrocks.http.url=192.168.xxx.xxx:8030
# Kafka トピック名が StarRocks テーブル名と異なる場合、それらの間のマッピング関係を設定する必要があります。
starrocks.topic2table.map=test:test_tbl
# StarRocks ユーザー名を入力します。
starrocks.username=user1
# StarRocks パスワードを入力します。
starrocks.password=123456
starrocks.database.name=example_db
sink.properties.strip_outer_array=trueNOTICE
ソースデータが Debezium フォーマットの CDC データであり、StarRocks テーブルが主キーテーブルである場合、ソースデータの変更を主キーテーブル に同期するために
transformを設定する必要があります。 -
Kafka Connect を設定して実行します。
-
Kafka Connect を設定します。config ディレクトリ内の設定ファイル config/connect-standalone.properties に以下のパラメータを設定します。詳細なパラメータと説明については、Running Kafka Connect を参照してください。以下の例では starrocks-kafka-connector バージョン
1.0.3を使用しています。新しいバージョンを使用する場合は、対応する変更を行う必要があります。# Kafka ブローカーのアドレス。複数の Kafka ブローカーのアドレスはカンマ (,) で区切る必要があります。
# この例では、Kafka クラスターにアクセスするためのセキュリティプロトコルとして PLAINTEXT を使用しています。他のセキュリティプロトコルを使用して Kafka クラスターにアクセスする場合は、このファイルに関連情報を設定する必要があります。
bootstrap.servers=<kafka_broker_ip>:9092
offset.storage.file.filename=/tmp/connect.offsets
offset.flush.interval.ms=10000
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=false
# 解凍後の starrocks-kafka-connector の絶対パス。例:
plugin.path=/home/kafka-connect/starrocks-kafka-connector-1.0.3 -
Kafka Connect を実行します。
CLASSPATH=/home/kafka-connect/starrocks-kafka-connector-1.0.3/* bin/connect-standalone.sh config/connect-standalone.properties config/connect-starrocks-sink.properties
-
分散モードで Kafka Connect を実行
-
Kafka Connect を設定して実行します。
-
Kafka Connect を設定します。config ディレクトリ内の設定ファイル
config/connect-distributed.propertiesに以下のパラメータを設定します。詳細なパラメータと説明については、Running Kafka Connect を参照してください。# Kafka ブローカーのアドレス。複数の Kafka ブローカーのアドレスはカンマ (,) で区切る必要があります。
# この例では、Kafka クラスターにアクセスするためのセキュリティプロトコルとして PLAINTEXT を使用しています。他のセキュリティプロトコルを使用して Kafka クラスターにアクセスする場合は、このファイルに関連情報を設定する必要があります。
bootstrap.servers=<kafka_broker_ip>:9092
offset.storage.file.filename=/tmp/connect.offsets
offset.flush.interval.ms=10000
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=false
# 解凍後の starrocks-kafka-connector の絶対パス。例:
plugin.path=/home/kafka-connect/starrocks-kafka-connector-1.0.3 -
Kafka Connect を実行します。
CLASSPATH=/home/kafka-connect/starrocks-kafka-connector-1.0.3/* bin/connect-distributed.sh config/connect-distributed.properties
-
-
Kafka コネクタを設定して作成します。分散モードでは、REST API を通じて Kafka コネクタを設定して作成する必要があります。パラメータと説明については、Parameters を参照してください。
注記Kafka コネクタはシンクコネクタです。
curl -i http://127.0.0.1:8083/connectors -H "Content-Type: application/json" -X POST -d '{
"name":"starrocks-kafka-connector",
"config":{
"connector.class":"com.starrocks.connector.kafka.StarRocksSinkConnector",
"topics":"test",
"key.converter":"org.apache.kafka.connect.json.JsonConverter",
"value.converter":"org.apache.kafka.connect.json.JsonConverter",
"key.converter.schemas.enable":"true",
"value.converter.schemas.enable":"false",
"starrocks.http.url":"192.168.xxx.xxx:8030",
"starrocks.topic2table.map":"test:test_tbl",
"starrocks.username":"user1",
"starrocks.password":"123456",
"starrocks.database.name":"example_db",
"sink.properties.strip_outer_array":"true"
}
}'備考ソースデータが Debezium フォーマットの CDC データであり、StarRocks テーブルが主キーテーブルである場合、ソースデータの変更を主キーテーブルに同期するために
transformを設定する必要があります。
StarRocks テーブルのクエリ
ターゲット StarRocks テーブル test_tbl をクエリします。
MySQL [example_db]> select * from test_tbl;
+------+-------------+
| id | city |
+------+-------------+
| 1 | New York |
| 2 | Los Angeles |
| 3 | Chicago |
+------+-------------+
3 rows in set (0.01 sec)
上記の結果が返された場合、データは正常にロードされています。
パラメータ
name
必須: YES
デフォルト値:
説明: この Kafka コネクタの名前。Kafka Connect クラスター内のすべての Kafka コネクタ間でグロー バルに一意である必要があります。例: starrocks-kafka-connector。
connector.class
必須: YES
デフォルト値:
説明: この Kafka コネクタのシンクで使用されるクラス。値を com.starrocks.connector.kafka.StarRocksSinkConnector に設定します。
topics
必須: YES
デフォルト値:
説明: 購読する1つ以上のトピックで、各トピックは StarRocks テーブルに対応します。デフォルトでは、StarRocks はトピック名が StarRocks テーブル名と一致すると仮定します。したがって、StarRocks はトピック名を使用してターゲットの StarRocks テーブルを決定します。topics または topics.regex (下記) のいずれかを選択して入力してください。ただし、StarRocks テーブル名がトピック名と異なる場合は、オプションの starrocks.topic2table.map パラメータ (下記) を使用してトピック名からテーブル名へのマッピングを指定します。
topics.regex
必須:
デフォルト値: 購読する1つ以上のトピックに一致する正規表現。詳細については topics を参照してください。topics.regex または topics (上記) のいずれかを選択して入力してください。
説明:
starrocks.topic2table.map
必須: NO
デフォルト値:
説明: トピック名が StarRocks テーブル名と異なる場合の StarRocks テーブル名とトピック名のマッピング。フォーマットは <topic-1>:<table-1>,<topic-2>:<table-2>,... です。
starrocks.http.url
必須: YES
デフォルト値:
説明: StarRocks クラスター内の FE の HTTP URL。フォーマットは <fe_host1>:<fe_http_port1>,<fe_host2>:<fe_http_port2>,... です。複数のアドレスはカンマ (,) で区切ります。例: 192.168.xxx.xxx:8030,192.168.xxx.xxx:8030。
starrocks.database.name
必須: YES
デフォルト値:
説明: StarRocks データベースの名前。
starrocks.username
必須: YES
デフォルト値:
説明: StarRocks クラスターアカウントのユーザー名。ユーザーは StarRocks テーブルに対する INSERT 権限を持っている必要があります。
starrocks.password
必須: YES
デフォルト値:
説明: StarRocks クラスターアカウントのパスワード。
key.converter
必須: NO
デフォルト値: Kafka Connect クラスターで使用されるキーコンバータ
説明: このパラメータは、シンクコネクタ (Kafka-connector-starrocks) のキーコンバータを指定し、Kafka データのキーをデシリアライズするために使用されます。デフォルトのキーコンバータは、Kafka Connect クラスターで使用されるものです。
value.converter
必須: NO
デフォルト値: Kafka Connect クラスターで使用される値コンバータ
説明: このパラメータは、シンクコネクタ (Kafka-connector-starrocks) の値コンバータを指定し、Kafka データの値をデシリアライズするために使用されます。デフォルトの値コンバータは、Kafka Connect クラスターで使用されるものです。
key.converter.schema.registry.url
必須: NO
デフォルト値:
説明: キーコンバータのスキーマレジストリ URL。
value.converter.schema.registry.url
必須: NO
デフォルト値:
説明: 値コンバータのスキーマレジストリ URL。
tasks.max
必須: NO
デフォルト値: 1
説明: Kafka コネクタが作成できるタスクスレッドの上限で、通常は Kafka Connect クラスターのワーカーノードの CPU コア数と同じです。このパラメータを調整してロードパフォーマンスを制御できます。
bufferflush.maxbytes
必須: NO
デフォルト値: 94371840(90M)
説明: 一度に StarRocks に送信される前にメモリに蓄積できるデータの最大サイズ。最大値は 64 MB から 10 GB の範囲です。Stream Load SDK バッファはデータをバッファリングするために複数の Stream Load ジョブを作成する可能性があることに注意してください。したがって、ここで言及されているしきい値は、総データサイズを指します。
bufferflush.intervalms
必須: NO
デフォルト値: 1000
説明: データのバッチを送信する間隔で、ロードの遅延を制御します。範囲: [1000, 3600000]。
connect.timeoutms
必須: NO
デフォルト値: 1000
説明: HTTP URL への接続のタイムアウト。範囲: [100, 60000]。
sink.properties.*
必須:
デフォルト値:
説明: ロード動作を制御するための Stream Load パラメータ。例えば、パラメータ sink.properties.format は Stream Load に使用されるフォーマット (CSV や JSON など) を指定します。サポートされているパラメータとその説明のリストについては、STREAM LOAD を参照してください。
sink.properties.format
必須: NO
デフォルト値: json
説明: Stream Load に使用されるフォーマット。Kafka コネクタは、データの各バッチをこのフォーマットに変換してから StarRocks に送信します。有効な値は csv と json です。詳細については、CSV パラメータ および JSON パラメータ を参照してください。
sink.properties.partial_update
必須: NO
デフォルト値: FALSE
説明: 部分更新を使用するかどうか。有効な値は TRUE と FALSE です。デフォルト値は FALSE で、この機能を無効にします。
sink.properties.partial_update_mode
必須: NO
デフォルト値: row
説明: 部分更新のモードを指定します。有効な値は row と column です。
- デフォルトの
rowは行モードでの部分更新を意味し、多くの列と小さなバッチでのリアルタイム更新に適しています。 columnは列モードでの部分更新を意味し、少ない列と多くの行でのバッチ更新に適しています。このようなシナリオでは、列モードを有効にすると更新速度が速くなります。例えば、100 列のテーブルで、すべての行に対して 10 列 (全体の 10%) のみが更新される場合、列モードの更新速度は 10 倍速くなります。
使用上の注意
フラッシュポリシー
Kafka コネクタはデータをメモリにバッファし、Stream Load を介して StarRocks にバッチでフラッシュします。以下の条件のいずれかが満たされた場合にフラッシュがトリガーされます。
- バッファされた行のバイト数が
bufferflush.maxbytesの制限に達したとき。 - 最後のフラッシュからの経過時間が
bufferflush.intervalmsの制限に達したとき。 - タスクのオフセットをコミットしようとするコネクタの間隔に達したとき。この間隔は Kafka Connect の設定
offset.flush.interval.msによって制御され、デフォルト値は60000です。
データの遅延を低くするために、これらの設定を Kafka コネクタの設定で調整します。ただし、フラッシュの頻度が増えると CPU と I/O の使用量が増加します。
制限
- Kafka トピックからの単一メッセージを複数のデータ行にフラット化して StarRocks にロードすることはサポートされていません。
- Kafka コネクタのシンクは少なくとも一度のセマンティクスを保証します。