すべてのプロダクト
Search
ドキュメントセンター

ApsaraDB for SelectDB:Kafka を使用したデータインポート

最終更新日:Aug 22, 2026

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:9092

Kafka の 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

コネクターのクラス名またはエイリアスです。この値は org.apache.doris.kafka.connector.DorisSinkConnector に設定します。

topics

ソースとなるトピックのコンマ区切りリストです。

doris.topic2table.map

トピックとテーブルのマッピングです。複数のマッピングはコンマ (,) で区切ります。例: topic1:tb1,topic2:tb2。このパラメーターが指定されていない場合、コネクターはトピック名とテーブル名が同じであると見なします。

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

データインポート方法です。次の方法がサポートされています:

  • stream_load:データを ApsaraDB for SelectDB に直接インポートします。

  • copy_into:データをオブジェクトストレージにインポートしてから ApsaraDB for SelectDB にロードします。

デフォルトは stream_load です。

sink.properties.*

Stream Load のインポートパラメーターです。

例: 列区切り文字を指定するには、sink.properties.column_separator=, を使用します。

詳細については、「Stream Load」をご参照ください。

delivery.guarantee

Kafka データを消費して ApsaraDB for SelectDB にインポートする際のデータ整合性の配信保証を指定します。サポートされているレベルは at_least_onceexactly_once で、デフォルトは at_least_once です。

現在、ApsaraDB for SelectDB では、copy into を使用したデータインポートでのみ exactly_once が保証されます。

enable.2pc

2 フェーズコミットを有効にして exactly-once セマンティクスを保証するかどうかを指定します。

説明

その他の一般的な Kafka Connect シンク設定については、「Configuring Connectors」をご参照ください。

前提条件

  1. バージョン 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.properties
  2. doris-kafka-connector-1.0.0.jar をダウンロードし、JAR ファイルを KAFKA_HOME/libs ディレクトリに配置します。

  3. ApsaraDB for SelectDB インスタンスを作成します。詳細については、「インスタンスの作成」をご参照ください。

  4. MySQL プロトコルを使用して ApsaraDB for SelectDB インスタンスに接続します。詳細については、「インスタンスへの接続」をご参照ください。

  5. テストデータベースとテストテーブルを作成します。

    1. テストデータベースを作成します。

      CREATE DATABASE test_db;
    2. テストテーブルを作成します。

      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 データの同期

  1. 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=1
  2. Kafka 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 を使用します。

  1. 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
  2. ダウンロードしたファイルを解凍します。

    tar -zxvf debezium-connector-mysql-1.9.8.Final-plugin.tar.gz
  3. 解凍したすべての JAR ファイルを KAFKA_HOME/libs ディレクトリに配置します。

  4. 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」をご参照ください。

  5. 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 にデータを同期する場合、事前にデータベースとテーブルを作成しておく必要があります。

  6. 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=test1234

SSL が有効な Kafka クラスターに接続するための Kafka Connect の設定に関する詳細については、「Configure Kafka Connect」をご参照ください。