ApsaraDB for SelectDB は、Doris Kafka Connector を使用して Kafka からデータを自動的にサブスクライブし、同期することをサポートしています。このトピックでは、Doris Kafka Connector を使用して ApsaraDB for SelectDB にデータを同期する方法について説明します。
背景情報
Kafka Connect は、Apache Kafka と他のシステム間でデータを確実にストリーミングするためのツールです。コネクターを定義して、大規模なデータセットを Kafka にインポートしたり、Kafka からエクスポートしたりできます。
Doris コミュニティが提供する Kafka コネクターは Kafka Connect クラスターで実行されます。Kafka トピックからデータを読み取り、ApsaraDB for SelectDB にデータを書き込みます。
ビジネスシナリオでは、通常、Debezium Connector を使用してデータベースの変更データを Kafka にプッシュするか、API を呼び出して JSON 形式のデータをリアルタイムで Kafka に書き込みます。Doris Kafka Connector は Kafka のデータを自動的にサブスクライブし、このデータを ApsaraDB for SelectDB に同期します。
Kafka Connect の実行モード
Kafka Connect には 2 つの実行モードがあります。
スタンドアロンモード
スタンドアロンモードは、本番環境では推奨されません。
スタンドアロンモードの設定
connect-standalone.properties ファイルを設定します。
# ブローカーアドレスを変更します
bootstrap.servers=127.0.0.1:9092Kafka の config ディレクトリに、connect-selectdb-sink.properties ファイルを作成し、次の内容を追加します。
name=test-selectdb-sink
connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
topics=topic_test
doris.topic2table.map=topic_test:test_kafka_tbl
buffer.count.records=10000
buffer.flush.time=120
buffer.size.bytes=5000000
doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com
doris.http.port=8030
doris.query.port=9030
doris.user=admin
doris.password=****
doris.database=test_db
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverterスタンドアロンモードでの起動
$KAFKA_HOME/bin/connect-standalone.sh -daemon $KAFKA_HOME/config/connect-standalone.properties $KAFKA_HOME/config/connect-selectdb-sink.properties分散モード
分散モードの設定
connect-distributed.properties ファイルを設定します。
# ブローカーアドレスを変更します
bootstrap.servers=127.0.0.1:9092
# group.id を変更します。ID は、同じクラスター内のすべてのワーカーで同じである必要があります。
group.id=connect-cluster分散モードでの起動
$KAFKA_HOME/bin/connect-distributed.sh -daemon $KAFKA_HOME/config/connect-distributed.propertiesコネクターの追加
curl -i http://127.0.0.1:8083/connectors -H "Content-Type: application/json" -X POST -d '{
"name":"test-selectdb-sink-cluster",
"config":{
"connector.class":"org.apache.doris.kafka.connector.DorisSinkConnector",
"topics":"topic_test",
"doris.topic2table.map": "topic_test:test_kafka_tbl",
"buffer.count.records":"10000",
"buffer.flush.time":"120",
"buffer.size.bytes":"5000000",
"doris.urls":"selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com",
"doris.user":"admin",
"doris.password":"***",
"doris.database":"test_db",
"doris.http.port":"8030",
"doris.query.port":"9030",
"key.converter":"org.apache.kafka.connect.storage.StringConverter",
"value.converter":"org.apache.kafka.connect.json.JsonConverter"
}
}'パラメーター
パラメーター | 説明 |
name | コネクターの名前です。ISO 制御文字を含まず、Kafka Connect 環境内で一意である必要があります。 |
connector.class | コネクターのクラス名またはエイリアスです。この値は |
topics | ソースとなるトピックのコンマ区切りリストです。 |
doris.topic2table.map | トピックとテーブルのマッピングです。複数のマッピングはコンマ (,) で区切ります。例: |
buffer.count.records | ApsaraDB for SelectDB にフラッシュされる前に、各 Kafka パーティションのメモリにバッファリングされるレコード数です。デフォルトは 10,000 です。 |
buffer.flush.time | メモリ内バッファーをフラッシュする間隔 (秒単位) です。デフォルトは 120 です。 |
buffer.size.bytes | 各 Kafka パーティションのメモリにバッファリングされるレコードの累積サイズ (バイト単位) です。デフォルトは 5,000,000 です。 |
doris.urls | ApsaraDB for SelectDB の接続エンドポイントです。 ApsaraDB for SelectDB コンソールの インスタンスの詳細 > ネットワーク情報 ページで関連パラメーターを取得できます。 例: selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com |
doris.http.port | ApsaraDB for SelectDB の HTTP ポートです。デフォルトは 8030 です。 |
doris.query.port | ApsaraDB for SelectDB の MySQL プロトコルポートです。デフォルトは 9030 です。 |
doris.user | ApsaraDB for SelectDB のユーザー名です。 |
doris.password | ApsaraDB for SelectDB のパスワードです。 |
doris.database | データが書き込まれる ApsaraDB for SelectDB データベースです。 |
key.converter | キー用の JSON コンバータークラスです。 |
value.converter | 値用の JSON コンバータークラスです。 |
jmx | JMX を介して内部コネクターメトリクスを取得するかどうかを指定します。詳細については、「Doris-Connector-JMX」をご参照ください。デフォルトは true です。 |
enable.delete | 削除操作を同期するかどうかを指定します。デフォルトは false です。 |
label.prefix | Stream Load を使用してインポートされたデータのラベルプレフィックスです。デフォルトはコネクターのアプリケーション名です。 |
auto.redirect | 有効にすると、コネクターはフロントエンド (FE) を介して Stream Load リクエストをターゲットのバックエンド (BE) にリダイレクトするため、BE 情報を取得する必要がなくなります。 |
load.model | データインポート方法です。次の方法がサポートされています:
デフォルトは |
sink.properties.* | Stream Load のインポートパラメーターです。 例: 列区切り文字を指定するには、 詳細については、「Stream Load」をご参照ください。 |
delivery.guarantee | Kafka データを消費して ApsaraDB for SelectDB にインポートする際のデータ整合性の配信保証を指定します。サポートされているレベルは 現在、ApsaraDB for SelectDB では、 |
enable.2pc | 2 フェーズコミットを有効にして exactly-once セマンティクスを保証するかどうかを指定します。 |
その他の一般的な Kafka Connect シンク設定については、「Configuring Connectors」をご参照ください。
例
前提条件
バージョン 2.4.0 以降の Apache Kafka クラスターまたは Confluent Cloud をインストールします。この例では、シングルノードの Kafka 環境を使用します。
# パッケージをダウンロードして解凍します wget https://archive.apache.org/dist/kafka/2.4.0/kafka_2.12-2.4.0.tgz tar -zxvf kafka_2.12-2.4.0.tgz cd kafka_2.12-2.4.0/ bin/zookeeper-server-start.sh -daemon config/zookeeper.properties bin/kafka-server-start.sh -daemon config/server.propertiesdoris-kafka-connector-1.0.0.jar をダウンロードし、JAR ファイルを KAFKA_HOME/libs ディレクトリに配置します。
ApsaraDB for SelectDB インスタンスを作成します。詳細については、「インスタンスの作成」をご参照ください。
MySQL プロトコルを使用して ApsaraDB for SelectDB インスタンスに接続します。詳細については、「インスタンスへの接続」をご参照ください。
テストデータベースとテストテーブルを作成します。
テストデータベースを作成します。
CREATE DATABASE test_db;テストテーブルを作成します。
USE test_db; CREATE TABLE employees ( emp_no int NOT NULL, birth_date date, first_name varchar(20), last_name varchar(20), gender char(2), hire_date date ) UNIQUE KEY(`emp_no`) DISTRIBUTED BY HASH(`emp_no`) BUCKETS 1;
例 1:JSON データの同期
SelectDB シンクの設定
スタンドアロンモードを例に、Kafka の config ディレクトリに selectdb-sink.properties ファイルを作成し、次の内容を追加します。
name=selectdb_sink connector.class=org.apache.doris.kafka.connector.DorisSinkConnector topics=test_topic doris.topic2table.map=test_topic:employees buffer.count.records=10000 buffer.flush.time=120 buffer.size.bytes=5000000 doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com doris.http.port=8030 doris.query.port=9030 doris.user=admin doris.password=*** doris.database=test_db key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.json.JsonConverter # オプション:デッドレターキューの設定 errors.tolerance=all errors.deadletterqueue.topic.name=test_error errors.deadletterqueue.context.headers.enable = true errors.deadletterqueue.topic.replication.factor=1Kafka Connect の起動
bin/connect-standalone.sh -daemon config/connect-standalone.properties config/selectdb-sink.properties
例 2:Debezium を使用した MySQL から ApsaraDB for SelectDB へのデータ同期
多くのビジネスシナリオでは、運用データベースからリアルタイムでデータを同期する必要があります。これには、データベースの変更データキャプチャ (CDC) メカニズムを使用する必要があります。
Debezium は Kafka Connect ベースの CDC ツールであり、MySQL、PostgreSQL、SQL Server、Oracle、MongoDB などのさまざまなデータベースに接続できます。データ変更を継続的に Kafka トピックに統一された形式で送信し、ダウンストリームのシンクがリアルタイムに消費します。この例では MySQL を使用します。
Debezium をダウンロードします。
wget https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/1.9.8.Final/debezium-connector-mysql-1.9.8.Final-plugin.tar.gzダウンロードしたファイルを解凍します。
tar -zxvf debezium-connector-mysql-1.9.8.Final-plugin.tar.gz解凍したすべての JAR ファイルを KAFKA_HOME/libs ディレクトリに配置します。
MySQL ソースを設定します。
Kafka の config ディレクトリに mysql-source.properties ファイルを作成し、次の内容を追加します。
name=mysql-source connector.class=io.debezium.connector.mysql.MySqlConnector database.hostname=rm-bp17372257wkz****.rwlb.rds.aliyuncs.com database.port=3306 database.user=testuser database.password=**** database.server.id=1 # Kafka におけるこのクライアントの一意の識別子 database.server.name=test123 # 同期するデータベースとテーブル。デフォルトでは、すべてのデータベースとテーブルが同期されます。 database.include.list=test table.include.list=test.test_table database.history.kafka.bootstrap.servers=localhost:9092 # データベーススキーマの変更を保存するために使用される Kafka トピック database.history.kafka.topic=dbhistory transforms=unwrap # https://debezium.io/documentation/reference/stable/transformations/event-flattening.html をご参照ください transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState # 削除イベントを記録 transforms.unwrap.delete.handling.mode=rewrite設定後、デフォルトの Kafka トピック名の形式は
SERVER_NAME.DATABASE_NAME.TABLE_NAMEになります。説明Debezium の設定については、「Debezium connector for MySQL」をご参照ください。
ApsaraDB for SelectDB シンクを設定します。
Kafka の config ディレクトリに selectdb-sink.properties ファイルを作成し、次の内容を追加します。
name=selectdb-sink connector.class=org.apache.doris.kafka.connector.DorisSinkConnector topics=test123.test.test_table doris.topic2table.map=test123.test.test_table:test_table buffer.count.records=10000 buffer.flush.time=120 buffer.size.bytes=5000000 doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com doris.http.port=8030 doris.query.port=9030 doris.user=admin doris.password=**** doris.database=test key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter # オプション:デッドレターキューの設定 #errors.tolerance=all #errors.deadletterqueue.topic.name=test_error #errors.deadletterqueue.context.headers.enable = true #errors.deadletterqueue.topic.replication.factor=1説明ApsaraDB for SelectDB にデータを同期する場合、事前にデータベースとテーブルを作成しておく必要があります。
Kafka Connect を起動します。
bin/connect-standalone.sh -daemon config/connect-standalone.properties config/mysql-source.properties config/selectdb-sink.properties説明起動後、
logs/connect.logファイルで、サービスが正常に起動したことを確認できます。
高度な使用方法
コネクターの操作
# コネクターのステータスを確認
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/status -X GET
# 現在のコネクターを削除
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster -X DELETE
# 現在のコネクターを一時停止
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/pause -X PUT
# 現在のコネクターを再開
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/resume -X PUT
# コネクター内のタスクを再起動
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/tasks/0/restart -X POST詳細については、「Connect REST Interface」をご参照ください。
デッドレターキュー
デフォルトでは、変換エラーが発生するとコネクターの処理は失敗します。ただし、エラーをスキップするようにコネクターを設定することで、このようなエラーを許容できます。また、エラーの詳細、失敗した操作、問題のあったレコードをデッドレターキューに書き込んで、後で分析することもできます。
errors.tolerance=all
errors.deadletterqueue.topic.name=test_error_topic
errors.deadletterqueue.context.headers.enable=true
errors.deadletterqueue.topic.replication.factor=1詳細については、「Error Reporting in Connect」をご参照ください。
SSL が有効な Kafka クラスターへの接続
Kafka Connect を介して SSL が有効な Kafka クラスターにアクセスするには、証明書ファイル (client.truststore.jks) を使用して Kafka ブローカーの公開キーを認証する必要があります。次の設定を connect-distributed.properties ファイルに追加できます。
# Connect ワーカー
security.protocol=SSL
ssl.truststore.location=/var/ssl/private/client.truststore.jks
ssl.truststore.password=test1234
# シンクコネクター用の組み込みコンシューマー
consumer.security.protocol=SSL
consumer.ssl.truststore.location=/var/ssl/private/client.truststore.jks
consumer.ssl.truststore.password=test1234SSL が有効な Kafka クラスターに接続するための Kafka Connect の設定に関する詳細については、「Configure Kafka Connect」をご参照ください。