SelectDB コネクタは、Realtime Compute for Apache Flink と ApsaraDB for SelectDB を統合します。ApsaraDB for SelectDB は、Alibaba Cloud 上のフルマネージドで Apache Doris と互換性のあるリアルタイムデータウェアハウスです。このコネクタを使用して、SelectDB のデータの読み取り、書き込み、またはルックアップを行うリアルタイムパイプラインを構築したり、YAML ベースのデータインジェストジョブでデータベース全体の同期を実行したりできます。
サポートされる機能:
| カテゴリ | 詳細 |
|---|---|
| テーブルタイプ | ソーステーブル、結果テーブル、ディメンションテーブル、データインジェストシンク |
| 実行モード | ストリームおよびバッチ |
| データフォーマット | JSON および CSV |
| API タイプ | DataStream、SQL、およびデータインジェスト YAML ジョブ |
| 更新/削除のサポート | はい |
| モニタリングメトリクス | なし |
主な特徴:
-
データベース全体のデータ同期
-
2フェーズコミット (2PC) による Exactly-once セマンティクス — レコードの重複や損失なし
-
Apache Doris 1.0 以降と互換
前提条件
開始する前に、以下を確認してください。
-
Ververica Runtime (VVR) 8.0.10 以降を搭載した Realtime Compute for Apache Flink
-
ApsaraDB for SelectDB インスタンス。詳細については、「インスタンスの作成」をご参照ください。
-
インスタンスに設定された IP アドレスホワイトリスト。詳細については、「ホワイトリストの設定」をご参照ください。
コネクタの設定
SelectDB コネクタは VVR 11.1 以降に組み込まれているため、手動でのインストールは不要です。
VVR 8.0.10 から 11.0 の場合は、コネクタを手動でインストールします。
-
Maven Central から JAR パッケージをダウンロードします (Flink バージョン 1.15–1.17)。
-
JAR を Realtime Compute for Apache Flink の開発コンソールにアップロードします。詳細については、「カスタムコネクタの管理」をご参照ください。
-
'connector' = 'doris'を使用して、SQL ジョブでコネクタを参照します。
SQL
構文
ソース、結果、ディメンションの3つのテーブルタイプはすべて同じ DDL 構文を共有します。含めるパラメーターによってテーブルのロールを指定します。
SelectDB をソーステーブルとして使用するには、まずクラスターへの直接接続を有効にする必要があります。ApsaraDB for SelectDB コンソールで、[インスタンス詳細] > [ネットワーク情報] に移動し、[クラスターへの直接接続を有効にする] をクリックします。これにより、高スループットの並列読み取りのために Arrow Flight SQL プロトコルが有効になります。
CREATE TABLE selectdb_source (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****'
);
パラメーター
全般
| パラメーター | 必須 | デフォルト | 説明 |
|---|---|---|---|
connector |
はい | — | doris に固定されます。 |
fenodes |
はい | — | SelectDB インスタンスの HTTP エンドポイント: <VPC アドレスまたはパブリックアドレス>:<HTTP プロトコルポート>。両方とも SelectDB コンソールの [インスタンス詳細] > [ネットワーク情報] から取得します。例: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080。 |
jdbc-url |
いいえ | — | ディメンションテーブルのルックアップとメタデータクエリ用の Java Database Connectivity (JDBC) 接続文字列: jdbc:mysql://<VPC アドレスまたはパブリックアドレス>:<MySQL プロトコルポート>。例: jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030。 |
table.identifier |
はい | — | <database>.<table> フォーマットのターゲットテーブル。例: db.tbl。 |
username |
はい | — | データベースのユーザー名。必要に応じて、[インスタンス詳細] ページの右上隅からパスワードをリセットします。 |
password |
はい | — | データベースユーザー名のパスワード。 |
doris.request.retries |
いいえ | 3 |
失敗したリクエストの再試行回数。 |
doris.request.connect.timeout |
いいえ | 30s |
接続タイムアウト。 |
doris.request.read.timeout |
いいえ | 30s |
読み取りタイムアウト。 |
ソーステーブル
| パラメーター | 必須 | デフォルト | 説明 |
|---|---|---|---|
doris.request.query.timeout |
いいえ | 21600s |
クエリタイムアウト (デフォルトは6時間)。 |
doris.request.tablet.size |
いいえ | 1 |
パーティションごとのタブレット数。値を小さくすると Flink の並列処理能力は向上しますが、データベースへの負荷が増加します。 |
doris.batch.size |
いいえ | 4064 |
リクエストごとにバックエンド (BE) ノードから読み取る最大行数。値を大きくすると、接続のオーバーヘッドとネットワーク遅延が減少します。 |
doris.exec.mem.limit |
いいえ | 8192mb |
クエリごとのメモリ制限 (バイト単位、デフォルトは 8 GB)。 |
source.use-flight-sql |
いいえ | false |
設定は不要です — SelectDB コンソールで [クラスターへの直接接続を有効にする] を有効にすると、Arrow Flight SQL が自動的に有効になります。 |
source.flight-sql-port |
いいえ | — | フロントエンド (FE) ノードの Arrow Flight SQL ポート (arrow_flight_sql_port)。 |
結果テーブル
書き込みモードは、配信保証とフラッシュ動作に影響します。一貫性の要件に基づいて選択してください。
| ストリーミング書き込み | バッチ書き込み | |
|---|---|---|
| トリガー条件 | Flink のチェックポイント間隔に従います | データ量または時間しきい値による定期的なフラッシュ |
| 配信保証 | Exactly-once (2PC 経由) | At-least-once (少なくとも1回)。Unique モデルでべき等性を実現します |
| レイテンシー | チェックポイント間隔に制約されます | 柔軟で、チェックポイントから独立しています |
| フォールトトレランス | Flink の完全な状態回復 | Unique モデルの重複排除に依存します |
| パラメーター | 必須 | デフォルト | 説明 |
|---|---|---|---|
sink.label-prefix |
いいえ | — | Stream Load インポートのラベルプレフィックス。すべてのジョブでグローバルに一意である必要があります — 同じラベルは一度しかコミットできません。ジョブの再起動をまたいで Exactly-once セマンティクスを保証するために必須です。 |
sink.properties.* |
いいえ | — | SelectDB Stream Load API に直接渡される Stream Load インポートパラメーター。以下の例をご参照ください。 |
sink.enable-delete |
いいえ | true |
DELETE 操作を伝播します。Doris テーブルでバッチ削除が有効になっている必要があり、Unique モデルでのみ機能します。 |
sink.enable-2pc |
いいえ | true |
Exactly-once セマンティクスのために2フェーズコミット (2PC) を有効にします。詳細については、「明示的なトランザクション操作」をご参照ください。 |
sink.buffer-size |
いいえ | 1 MB |
書き込みキャッシュバッファーのサイズ (バイト単位)。デフォルトのままにしてください。 |
sink.buffer-count |
いいえ | 3 |
書き込みキャッシュバッファーの数。デフォルトのままにしてください。 |
sink.max-retries |
いいえ | 3 |
コミット失敗後の最大再試行回数。 |
sink.enable.batch-mode |
いいえ | false |
バッチ書き込みモードに切り替えます。フラッシュはチェックポイントではなく、以下の3つの sink.buffer-flush.* パラメーターによって制御されます。Exactly-once は保証されません。べき等性のために Unique モデルを使用してください。 |
sink.flush.queue-size |
いいえ | 2 |
バッチモードでのキャッシュキューのサイズ。 |
sink.buffer-flush.max-rows |
いいえ | 500000 |
バッチモードでのフラッシュごとの最大行数。 |
sink.buffer-flush.max-bytes |
いいえ | 100 MB |
バッチモードでのフラッシュごとの最大バイト数。 |
sink.buffer-flush.interval |
いいえ | 10s |
バッチモードでのフラッシュ間隔。 |
sink.ignore.update-before |
いいえ | true |
Flink CDC からの update-before イベントを無視します。 |
**sink.properties.* の例:**
CSV フォーマット:
'sink.properties.column_separator' = ','
-- 値にカンマが含まれる可能性がある場合は、印刷不可能な区切り文字を使用します:
-- 'sink.properties.column_separator' = '\x01'
JSON フォーマット:
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true'
-- または: 'sink.properties.strip_outer_array' = 'true'
ディメンションテーブル
| パラメーター | 必須 | デフォルト | 説明 |
|---|---|---|---|
lookup.cache.max-rows |
いいえ | -1 |
ルックアップキャッシュの最大行数。-1 はキャッシュを無効にします。 |
lookup.cache.ttl |
いいえ | 10s |
キャッシュエントリの生存時間 (TTL)。 |
lookup.max-retries |
いいえ | 1 |
ルックアップ クエリ失敗時の再試行。 |
lookup.jdbc.async |
いいえ | false |
非同期ルックアップを有効にします。 |
lookup.jdbc.read.batch.size |
いいえ | 128 |
非同期ルックアップモードでのクエリごとの最大バッチサイズ。 |
lookup.jdbc.read.batch.queue-size |
いいえ | 256 |
非同期ルックアップモードでの中間バッファーキューのサイズ。 |
lookup.jdbc.read.thread-size |
いいえ | 3 |
非同期ルックアップモードでのタスクごとの JDBC ルックアップスレッド数。 |
例
ソーステーブル
CREATE TEMPORARY TABLE selectdb_source (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****'
);
結果テーブル
CREATE TEMPORARY TABLE selectdb_sink (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****',
'sink.label-prefix' = 'flink_orders' -- ジョブ間でグローバルに一意である必要があります
);
ディメンションテーブル
SelectDB は、ストリーミングファクトテーブルに対して結合されるルックアップディメンションテーブルとして機能します。
-- Kafka からのファクトテーブル
CREATE TEMPORARY TABLE fact_table (
`id` BIGINT,
`name` STRING,
`city` STRING,
`process_time` AS proctime()
) WITH (
'connector' = 'kafka',
...
);
-- SelectDB からのディメンションテーブル
CREATE TEMPORARY TABLE dim_city (
`city` STRING,
`level` INT,
`province` STRING,
`country` STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'jdbc-url' = 'jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030',
'table.identifier' = 'dim.dim_city',
'username' = 'admin',
'password' = '****'
);
-- テンポラル結合
SELECT a.id, a.name, a.city, c.province, c.country, c.level
FROM fact_table a
LEFT JOIN dim_city FOR SYSTEM_TIME AS OF a.process_time AS c
ON a.city = c.city;
データインジェスト
データベース全体の同期のために、YAML ベースのデータインジェストジョブで SelectDB コネクタをシンクとして使用します。
構文
source:
type: <source-type>
sink:
type: doris
name: Doris Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
username: root
password: ""
パラメーター
| パラメーター | 必須 | デフォルト | 説明 |
|---|---|---|---|
type |
はい | — | doris に固定されます。 |
name |
いいえ | — | シンクの説明的な名前。 |
fenodes |
はい | — | HTTP エンドポイント: <VPC アドレスまたはパブリックアドレス>:<HTTP プロトコルポート>。両方とも SelectDB コンソールの [インスタンス詳細] > [ネットワーク情報] から取得します。例: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080。 |
jdbc-url |
いいえ | — | JDBC 接続文字列。例: jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030。 |
username |
はい | — | データベースのユーザー名。 |
password |
はい | — | データベースユーザー名のパスワード。 |
sink.enable.batch-mode |
いいえ | true |
データインジェストジョブでは、バッチモードがデフォルトでオンになっています。フラッシュは3つの sink.buffer-flush.* パラメーターによって制御されます。Exactly-once は保証されません。べき等性のために Unique モデルを使用してください。 |
sink.flush.queue-size |
いいえ | 2 |
キャッシュキューのサイズ。 |
sink.buffer-flush.max-rows |
いいえ | 500000 |
フラッシュごとの最大行数。 |
sink.buffer-flush.max-bytes |
いいえ | 100 MB |
フラッシュごとの最大バイト数。 |
sink.buffer-flush.interval |
いいえ | 10s |
フラッシュ間隔。最小: 1s。 |
sink.properties.* |
いいえ | — | Stream Load インポートパラメーター。 |
**sink.properties.* の例:**
CSV フォーマット:
sink.properties.column_separator: ','
# 値にカンマが含まれる可能性がある場合は、印刷不可能な区切り文字を使用します:
# sink.properties.column_separator: '\x01'
JSON フォーマット:
sink.properties.format: 'json'
sink.properties.read_json_by_line: 'true'
例
以下の例は、一般的なデータインジェストのシナリオを対象としています。プレースホルダーのエンドポイント、ユーザー名、パスワードを実際の値に置き換えてください。
${secret_values.variable_name} は、ワークスペースで作成されたシークレット変数を参照します。詳細については、「変数の管理」をご参照ください。ソースパラメーターの完全な説明については、「MySQL YAML コネクタ」および「Postgres CDC コネクタ」をご参照ください。
コネクタは宛先データベースを作成しません。開始する前に、以下の例で使用されている宛先データベースを ApsaraDB for SelectDB インスタンスに作成してください。宛先テーブルは存在しない場合、自動的に作成されます。
単一テーブルの同期
単一の MySQL テーブルを SelectDB に同期します。宛先テーブルは存在しない場合、自動的に作成されます。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
# (任意) 最初に既存の完全データを同期し、次に増分データを同期します。
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
pipeline:
name: MySQL to SelectDB Pipeline
データベース全体の同期
MySQL データベース内のすべてのテーブルを SelectDB に同期します。ダウンストリームのデータベース名とテーブル名は、アップストリームのものと同じです。宛先テーブルは存在しない場合、自動的に作成されます。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
# test_db データベース内のすべてのテーブルに一致します。
tables: test_db.\.*
server-id: 5401-5499
scan.startup.mode: initial
# (任意) ジョブを再起動せずに、増分フェーズ中に作成されたテーブルを同期します。
scan.binlog.newly-added-table.enabled: true
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
# (任意) バッチ書き込みのフラッシュ間隔。小規模なワークロードでは、データがバッファーに長時間保持されないように、この値を小さくします。
sink.buffer-flush.interval: 10s
pipeline:
name: MySQL to SelectDB Pipeline
データベース全体の同期からテーブルを除外
一時テーブルやログテーブルなど、同期したくないテーブルをスキップします。ルートモジュールがない場合、データはソーステーブルと同じ名前の宛先テーブルに書き込まれます。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.\.*
# この正規表現に一致するテーブルは同期されません。複数の正規表現はカンマ (,) で区切ります。
tables.exclude: test_db.tmp_\.*, test_db.log_\.*
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
pipeline:
name: MySQL to SelectDB Pipeline
指定されたデータベースとテーブルへの同期
ダウンストリームのデータベース名またはテーブル名がアップストリームのものと異なる場合は、ルートモジュールを使用して宛先の名前を変更します。次の例では、test_db のテーブルを ods_db データベースに同期し、各テーブル名に ods_ プレフィックスを追加します。たとえば、test_db.orders は ods_db.ods_orders に同期されます。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.\.*
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
route:
# <> は、一致したソーステーブル名に置き換えられるプレースホルダーです。
- source-table: test_db.\.*
sink-table: ods_db.ods_<>
replace-symbol: <>
pipeline:
name: MySQL to SelectDB Pipeline
シャードテーブルのマージ
同じスキーマを共有する複数のシャードテーブルを単一の SelectDB テーブルにマージします。次の例では、変換モジュールを使用してソースデータベースとテーブルを識別する列を追加し、この列とビジネスプライマリキーを複合プライマリキーとして使用して、マージされたデータの出所を追跡できるようにします。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
# order_db_1.orders_1 や order_db_2.orders_2 など、すべてのシャードに一致します。
tables: order_db_\d+.orders_\d+
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
# ダウンストリームのプライマリキーがアップストリームのものと異なる場合は、このパラメーターを false に設定します。そうしないと、古いプライマリキーに対応する行が削除されません。
sink.ignore.update-before: false
transform:
- source-table: order_db_\d+.orders_\d+
# ソースデータベースとテーブルを識別する src_table 列を追加します。
projection: "*, __schema_name__ || '.' || __table_name__ AS src_table"
# 複合プライマリキーは、宛先テーブルのキー列と同じでなければなりません。
primary-keys: order_id, src_table
description: ソース識別子を追加し、複合プライマリキーを設定します
route:
# すべてのシャードを単一のテーブルにマージします。
- source-table: order_db_\d+.orders_\d+
sink-table: dw_db.merged_orders
pipeline:
name: MySQL sharding to SelectDB Pipeline
データのフィルタリングと列のプルーニング
変換モジュールを使用して、データをフィルタリングし、列をプルーニングします。フィルター条件はソーステーブルのフィールドで評価され、レコードを同期するかどうかを決定し、プロジェクションはダウンストリームテーブルに書き込む列を決定します。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
transform:
- source-table: test_db.test_source_table
# order_status が PAID のレコードのみを同期します。
filter: "order_status = 'PAID'"
# 以下の4つの列のみを書き込みます。
projection: order_id, customer_id, total_amount, created_at
primary-keys: order_id
description: 注文ステータスでフィルタリングし、列をプルーニングします
pipeline:
name: MySQL to SelectDB Pipeline
PostgreSQL データベース全体の同期
PostgreSQL スキーマ内のすべてのテーブルを SelectDB に同期します。開始する前に、PostgreSQL インスタンスで wal_level を logical に設定し、max_replication_slots と max_wal_senders に十分な余裕があることを確認してください。詳細については、「PostgreSQL データベースの設定」をご参照ください。
source:
type: postgres
name: PostgreSQL Source
hostname: <yourHostname>
port: 5432
username: ${secret_values.pg_username}
password: ${secret_values.pg_password}
# PostgreSQL のテーブル名は database.schema.table フォーマットを使用します。
tables: test_db.public.\.*
slot.name: <yourSlotName>
decoding.plugin.name: pgoutput
scan.startup.mode: initial
# (任意) ジョブ開始時にパブリケーションを作成する代わりに、既存のパブリケーションを使用します。
debezium.publication.autocreate.mode: disabled
debezium.publication.name: <yourPublicationName>
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
pipeline:
name: PostgreSQL to SelectDB Pipeline
型マッピング
Flink から SelectDB へ
| Flink CDC 型 | SelectDB 型 | 注記 |
|---|---|---|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INT |
INT |
|
BIGINT |
BIGINT |
|
DECIMAL |
DECIMAL |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
BOOLEAN |
BOOLEAN |
|
DATE |
DATE |
|
TIMESTAMP[(p)] |
DATETIME[(p)] |
|
TIMESTAMP_LTZ[(p)] |
DATETIME[(p)] |
|
CHAR(n) |
CHAR(n*3) |
SelectDB は文字列を UTF-8 で格納します。英字は 1 バイト、漢字は 3 バイトを占有します。最大 CHAR 長は 255 です。それより長い値は自動的に VARCHAR に変換されます。 |
VARCHAR(n) |
VARCHAR(n*3) |
同じ UTF-8 の乗数が適用されます。最大 VARCHAR 長は 65533 です。それより長い値は自動的に STRING に変換されます。 |
BINARY(n) |
STRING |
|
VARBINARY(n) |
STRING |
|
STRING |
STRING |