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

Realtime Compute for Apache Flink:SelectDB

最終更新日:Sep 16, 2026

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 の場合は、コネクタを手動でインストールします。

  1. Maven Central から JAR パッケージをダウンロードします (Flink バージョン 1.15–1.17)。

  2. JAR を Realtime Compute for Apache Flink の開発コンソールにアップロードします。詳細については、「カスタムコネクタの管理」をご参照ください。

  3. '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

次のステップ