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

Realtime Compute for Apache Flink:OceanBase

最終更新日:Sep 03, 2026

OceanBase コネクタは、Realtime Compute for Apache Flink と、ネイティブな分散型ハイブリッドトランザクション/分析処理 (HTAP) データベースである OceanBase を統合します。このコネクタを使用して、変更データキャプチャ (CDC) ストリームを読み取り、ディメンションテーブルを結合し、結果を OceanBase に書き込むことができます。

OceanBase コネクタはパブリックプレビュー中です。

背景情報

OceanBase は、ネイティブ分散型のハイブリッド トランザクションおよび分析処理 (HTAP) データベース管理システムです。詳細については、OceanBase の公式サイトをご参照ください。MySQL や Oracle データベースから移行する際の業務システムのリファクタリングコストを削減するため、OceanBase は Oracle モードと MySQL モードの 2 つの互換性モードをサポートしています。これらのモードでは、データ型、SQL 機能、および内部ビューが、MySQL や Oracle のものと一致しています。

サポート機能

カテゴリ 詳細
テーブルタイプ ソース、ディメンション、シンクテーブル
ランタイムモード ストリーミングモードとバッチモード
データフォーマット 該当なし
特定の監視メトリクス なし
API SQL
シンクテーブルでの更新と削除 はい

前提条件

開始する前に、次のことを確認してください:

制限

  • Ververica Runtime (VVR) 8.0.1 以降が必要です。

セマンティクス保証:

  • CDC ソーステーブル: exactly-once セマンティクス。障害発生後でも、完全な履歴データから Binlog の読み取りに移行する際に、データの損失や重複は発生しません。

  • シンクテーブル: at-least-once セマンティクス。シンクテーブルにプライマリキーがある場合、べき等性がデータの正確性を保証します。

VVR 11.4.0 における CDC アーキテクチャ変更

VVR 11.4.0 以降、OceanBase CDC コネクタは次のとおりアップグレードされました:

  • OceanBase LogProxy サービスに基づく従来の CDC コネクタは、非推奨となり、削除されました。

  • 増分ログキャプチャには、OceanBase Binlog サービス が必要になりました。OceanBase CDC コネクタは、標準の MySQL CDC コネクタを直接接続する場合と比べて、Binlog サービスとのプロトコル互換性と接続安定性が優れています。変更追跡のために、標準の MySQL CDC コネクタを OceanBase Binlog サービスに接続することは推奨されません。

  • Oracle 互換モードでの増分変更の追跡はサポートされなくなりました。Oracle モードの CDC については、OceanBase Enterprise テクニカルサポートにお問い合わせください。

構文

CREATE TABLE oceanbase_source (
   order_id     INT,
   order_date   TIMESTAMP(0),
   customer_name STRING,
   price        DECIMAL(10, 5),
   product_id   INT,
   order_status BOOLEAN,
   PRIMARY KEY(order_id) NOT ENFORCED
) WITH (
  'connector'  = 'oceanbase',
  'url'        = '<your-jdbc-url>',
  'tableName'  = '<your-table-name>',
  'userName'   = '<your-username>',
  'password'   = '<your-password>'
);

シンクの書き込み動作:

コネクタは、受信するレコードごとに、シンクテーブルのスキーマに基づいて SQL ステートメントを生成します:

  • プライマリキーなし: INSERT INTO

  • プライマリキーあり: UPSERT (データベースの互換性に応じて)

WITH パラメーター

一般パラメーター

これらのパラメーターは、すべてのテーブルタイプに適用されます。

パラメーター 説明 必須 タイプ デフォルト
connector oceanbase に設定します。 はい 文字列 —
password データベースのパスワード。 はい 文字列 —

ソーステーブルのパラメーター

重要

VVR 11.4.0 以降、OceanBase CDC コネクターは増分ログキャプチャに OceanBase Binlog サービスを使用します。LogProxy ベースのコネクターは削除されました。VVR 11.4.0 から Oracle 互換モードの CDC はサポートされなくなりました。

パラメーター 説明 必須 タイプ デフォルト 注意
hostname OceanBase データベースの IP アドレスまたはホスト名。可能な場合は、Virtual Private Cloud (VPC) アドレスを使用してください。 はい STRING — OceanBase と Realtime Compute for Apache Flink が異なる VPC にある場合は、VPC 間接続を設定するか、パブリックエンドポイントを使用してください。詳細については、「Workspace management」および「Internet access for Flink clusters」をご参照ください。
username OceanBase データベースのユーザー名。 はい STRING — —
database-name OceanBase データベース名。正規表現をサポートし、複数のデータベースから読み取ることができます。^ および $ アンカーは使用しないでください。 はい STRING — コネクターは、database-name と table-name を \\. (VVR 8.0.1 以降) または . (それ以前のバージョン) で連結して、完全なパスの正規表現を作成します。例: db_.* + tb_.+ は db_.*\\.tb_.+ になります。
table-name OceanBase テーブル名。正規表現をサポートし、複数のテーブルから読み取ることができます。^ および $ アンカーは使用しないでください。 はい STRING — 上記の database-name の注意をご参照ください。
port OceanBase データベースのポート。 いいえ INTEGER 3306 —
server-id データベースクライアントの数値 ID。グローバルに一意である必要があります。5400-5408 のような範囲をサポートし、同時リーダーに異なる ID を割り当てることができます。 いいえ STRING 5400 から 6400 までのランダムな値 同じデータベースに接続する各ジョブには、異なる ID を使用してください。詳細については、「Server ID usage」をご参照ください。
scan.incremental.snapshot.chunk.size 増分スナップショットの読み取り中におけるチャンクあたりの行数。各チャンクのデータは、完全に読み取られる前にメモリにバッファリングされます。チャンクを小さくすると障害復旧の粒度が向上しますが、メモリ不足 (OOM) エラーが発生したり、スループットが低下したりする可能性があります。 いいえ INTEGER 8096 チャンクサイズと、メモリおよびスループットの要件とのバランスを取ってください。
scan.snapshot.fetch.size テーブルのフル読み取り中に、プルごとにフェッチされるレコードの最大数。 いいえ INTEGER 1024 —
scan.startup.mode データ消費の起動モード。 いいえ STRING initial 有効な値: initial (完全な履歴をスキャンしてから Binlog を読み取る)、 latest-offset (Binlog の末尾からのみ読み取り)、 earliest-offset (利用可能な最も古い Binlog)、 specific-offset (scan.startup.specific-offset.* パラメーターで指定)、 timestamp (scan.startup.timestamp-millis で指定)。
重要

earliest-offset、specific-offset、および timestamp モードの場合、指定された Binlog の位置からジョブの起動までの間にテーブルスキーマが変更されていない必要があります。

scan.startup.specific-offset.file 開始オフセットの Binlog ファイル名。例: mysql-bin.000003。 いいえ STRING — scan.startup.mode=specific-offset が必要です。
scan.startup.specific-offset.pos 指定された Binlog ファイル内のバイトオフセット。 いいえ INTEGER — scan.startup.mode=specific-offset が必要です。
scan.startup.specific-offset.gtid-set 開始オフセットの GTID セット。例: 24DA167-0C0C-11E8-8442-00059A3C7B00:1-19。 いいえ STRING — scan.startup.mode=specific-offset が必要です。
scan.startup.timestamp-millis 開始タイムスタンプ (ミリ秒単位)。OceanBase CDC は、各 Binlog ファイルの初期イベントを読み取り、このタイムスタンプに一致するファイルを見つけます。Binlog ファイルはパージされていない必要があります。 いいえ LONG — scan.startup.mode=timestamp が必要です。
server-time-zone データベースが使用するセッションタイムゾーン。TIMESTAMP 型が STRING に変換される方法を制御します。詳細については、「Debezium temporal types」をご参照ください。例: Asia/Shanghai。 いいえ STRING Flink ジョブのランタイムタイムゾーン —
debezium.min.row.count.to.stream.results コネクターがフル読み取り (テーブル全体をメモリに) からバッチ読み取り (行をバッチでストリーミング) に切り替える行数のしきい値。フル読み取りの方が高速です。バッチ読み取りは、大きなテーブルでの OOM を回避します。 いいえ INTEGER 1000 —
connect.timeout 接続がタイムアウトしてから再試行するまでの最大待機時間。 いいえ DURATION 30s —
connect.max-retries 障害発生後の接続再試行の最大数。 いいえ INTEGER 3 —
connection.pool.size データベースのコネクションプールのサイズ。接続を再利用すると、開いている接続の総数が減少します。 いいえ INTEGER 20 —
jdbc.properties.* カスタム JDBC URL 接続パラメーター。例: 'jdbc.properties.useSSL' = 'false'。詳細については、「MySQL configuration properties」をご参照ください。 いいえ STRING — —
debezium.* Binlog 読み取り用のカスタム Debezium パラメーター。例: 'debezium.event.deserialization.failure.handling.mode' = 'ignore'。 いいえ STRING — —
heartbeat.interval ソースがハートビートイベントを発行して Binlog オフセットを進める間隔。更新頻度の低いテーブルで Binlog オフセットが期限切れになるのを防ぎます。オフセットが期限切れになるとジョブは失敗し、ステートレスな再起動が必要になります。 いいえ DURATION 30s —
scan.incremental.snapshot.chunk.key-column スナップショットフェーズ中にデータを分割するためのチャンクキーとして使用される列。 条件付き STRING — プライマリキーがないテーブルの場合に必須です (NOT NULL である必要があります)。プライマリキーのあるテーブルではオプションです (プライマリキーから 1 つの列を選択します)。
scan.incremental.close-idle-reader.enabled スナップショットフェーズの完了後にアイドル状態のリーダーを閉じるかどうかを指定します。 いいえ BOOLEAN false VVR 8.0.1 以降。execution.checkpointing.checkpoints-after-tasks-finish.enabled=true も必要です。
scan.read-changelog-as-append-only.enabled 変更ログストリームを追加専用ストリームに変換するかどうかを指定します。true の場合、すべてのメッセージタイプ (INSERT、DELETE、UPDATE_BEFORE、UPDATE_AFTER) が INSERT に変換されます。アップストリームテーブルからの削除メッセージを保持する場合など、特別なシナリオでのみ有効にしてください。 いいえ BOOLEAN false VVR 8.0.8 以降。
scan.only.deserialize.captured.tables.changelog.enabled 増分フェーズ中に、キャプチャされたテーブルの変更イベントのみを逆シリアル化するかどうかを指定します。true に設定すると、Binlog の読み取りが高速化されます。 いいえ BOOLEAN false (VVR 8.x)、true (VVR 11.1 以降) VVR 8.0.7 以降。VVR 8.0.8 以前では、パラメーター名 debezium.scan.only.deserialize.captured.tables.changelog.enable を使用してください。
scan.parse.online.schema.changes.enabled 増分フェーズ中に、ApsaraDB RDS のロックレスな変更の DDL イベントを解析するかどうかを指定します。実験的な機能です。オンラインでロックレスなスキーマ変更を実行する前に、Flink ジョブのスナップショットを取得してください。 いいえ BOOLEAN false VVR 11.1 以降。
scan.incremental.snapshot.backfill.skip スナップショットフェーズ中にバックフィルをスキップするかどうかを指定します。バックフィルは、単一チャンクのスナップショットクエリ中にのみ適用され、フル読み取りフェーズ全体をカバーするものではありません。バックフィルがスキップされると、各チャンクのスナップショットクエリはその時点での最新のテーブルデータを読み取ります。チャンクが読み取られた後に発生した更新は、フル読み取りフェーズ中にはマージされず、増分フェーズに入った後に Binlog から読み取られます。例: chunk5 のスナップショット作成中に chunk5 に更新が発生した場合、その更新は chunk5 のスナップショットに直接反映されます。リーダーが chunk80 に進んだ後に chunk5 が更新された場合、その更新は後で増分フェーズ中に Binlog から適用されます。重要: この機能を有効にすると、チャンクのスキャン中またはスキャン後に発生した変更は、増分フェーズで Binlog から配信されるため、重複する可能性があります。at-least-once セマンティクスのみが保証されます。この機能は、ダウンストリームシンクがプライマリキーによるべき等な書き込みをサポートしている場合にのみ有効にしてください。 いいえ BOOLEAN false VVR 11.1 以降。
scan.incremental.snapshot.unbounded-chunk-first.enabled スナップショットフェーズ中に、境界のないチャンクを最初に分散するかどうかを指定します。TaskManager が最後のチャンクを処理する際の OOM のリスクを軽減します。実験的な機能です。ジョブを初めて開始する前に、このパラメーターを追加してください。 いいえ BOOLEAN false VVR 11.1 以降。

ディメンションテーブルパラメーター

パラメーター 説明 必須 タイプ デフォルト 備考
url JDBC URL。MySQL データベース名または Oracle サービス名を含める必要があります。 はい STRING — —
userName データベースのユーザー名。 はい STRING — —
cache ディメンションテーブルルックアップのキャッシュポリシー。 いいえ STRING ALL ALL :ジョブ開始前にすべてのデータをロードし、有効期限が切れた後に再読み込みします。ルックアップミスが多い小さなテーブルに適しています。非同期ロードをサポートするために、結合ノードのメモリをテーブルサイズの少なくとも 2 倍に増やしてください。LRU :行のサブセットをキャッシュします。cacheSize が必要です。None :キャッシュなし。
cacheSize キャッシュされたエントリの最大数。 いいえ INTEGER 100000 cache=LRU の場合は必須です。cache=ALL の場合は無視されます。
cacheTTLMs キャッシュタイムアウト (ミリ秒単位)。動作は cache の設定によって異なります。LRU の場合、エントリはこの期間後に有効期限が切れます (デフォルトでは有効期限なし)。ALL の場合、この期間後にフルキャッシュが再読み込みされます (デフォルトでは再読み込みなし)。None の場合、このパラメーターは効果がありません。 いいえ LONG Long.MAX_VALUE —
maxRetryTimeout 最大リトライ期間。 いいえ DURATION 60s —

シンクテーブルパラメーター (JDBC)

パラメーター 説明 必須 型 デフォルト 備考
url JDBC URL です。MySQL のデータベース名または Oracle のサービス名を含める必要があります。 はい 文字列 — —
userName ユーザー名。 はい 文字列 — —
tableName ターゲットテーブル名。 はい 文字列 — —
sink.mode 書き込みモード。標準書き込みの場合は jdbc に設定します。バイパスインポートの場合は direct-load に設定します。 はい 文字列 jdbc —
compatibleMode OceanBase の互換モード。有効な値: mysql、oracle。 いいえ 文字列 mysql OceanBase 固有のパラメーター。
maxRetryTimes 書き込みの最大再試行回数。 いいえ 整数 3 —
poolInitialSize 初期コネクションプールサイズ。 いいえ 整数 1 —
poolMaxActive プール内のアクティブな接続の最大数。 いいえ 整数 8 —
poolMaxWait プールから接続を取得するための最大待機時間 (ミリ秒)。 いいえ 整数 2000 —
poolMinIdle プール内のアイドル接続の最小数。 いいえ 整数 1 —
connectionProperties k1=v1;k2=v2 形式の JDBC 接続プロパティ。 いいえ 文字列 — —
ignoreDelete 削除操作を無視するかどうかを指定します。 いいえ ブール値 false —
excludeUpdateColumns 更新から除外する列をカンマ区切りで指定します (例:column1,column2)。プライマリキー列は、この設定に関係なく常に除外されます。 いいえ 文字列 — —
partitionKey パーティションキー。設定すると、modRule が適用される前に、このキーでデータがグループ化されます。 いいえ 文字列 — —
modRule column_name mod number 形式のグループ化ルール (例:user_id mod 8)。列は数値型である必要があります。データは最初に partitionKey でパーティション化され、次に各パーティション内でこのルールによってグループ化されます。 いいえ 文字列 — —
bufferSize データバッファーのサイズ (レコード数)。 いいえ 整数 1000 —
flushIntervalMs バッファーのフラッシュ間隔 (ミリ秒)。この間隔内にバッファーが出力条件を満たさない場合、バッファー内のすべてのデータが自動的にフラッシュされます。 いいえ 長整数 1000 —
retryIntervalMs 再試行間隔 (ミリ秒)。 いいえ 整数 5000 —

シンクテーブルのパラメータ(バイパスインポート)

バイパスインポートは、OceanBase でのバルクデータロードのための高スループットな書き込み方法です。VVR 11.5 以降で利用できます。

バイパスインポートを使用する前に、以下の制約を確認してください。

  • 有界ストリームのみ:データソースは有界ストリームである必要があります。最高のパフォーマンスを得るには、Flink バッチモードを使用してください。

  • インポート中のテーブルロック:ターゲットテーブルは、インポート全体を通じてロックされます。DML 書き込みと DDL 変更はブロックされますが、読み取りクエリは影響を受けません。

  • リアルタイム書き込みには不適:ストリーミングまたはリアルタイム書き込みの場合は、代わりに JDBC シンクを使用してください。

パラメーター 説明 必須 型 デフォルト 備考
sink.mode バイパスインポートを使用するには、 direct-load に設定します。 いいえ 文字列 jdbc —
host OceanBase データベースの IP アドレスまたはホスト名。 はい 文字列 — —
port OceanBase データベースの RPC ポート。 いいえ 整数 2882 —
username データベースのユーザー名。 はい 文字列 — —
tenant-name OceanBase のテナント名。 はい 文字列 — —
schema-name MySQL テナントの場合:データベース名。Oracle テナントの場合:所有者名。 はい 文字列 — —
table-name ターゲットテーブル名。 はい 文字列 — —
parallel インポートタスクのサーバーサイドの同時実行数。サーバーは、エラーを返さずに、テナントの CPU 仕様に基づいて実際の並列度を制限します。計算式: MIN(tenant_cores × 2, parallel) × partition_nodes。たとえば、2 CPU コア、 parallel=10、2 パーティションノードの場合: MIN(4, 10) × 2 = 8。 いいえ 整数 8 —
buffer-size OceanBase への 1 回の書き込み前にバッファリングされるレコード数。 いいえ 整数 1024 —
dup-action 重複するプライマリキーが検出された場合の動作。 STOP_ON_DUP:インポートが失敗します。 REPLACE:既存の行を上書きします。 IGNORE:受信した行を破棄します。 いいえ 文字列 REPLACE —
load-method インポートモード。 full:標準バイパスインポート。 inc:増分モード、プライマリキーの競合をチェックします (observer 4.3.2 以降、 dup-action=REPLACE はサポートされていません)。 inc_replace:増分置換モード、競合チェックなしで既存の行を直接上書きします (observer 4.3.2 以降、 dup-action は無視されます)。 いいえ 文字列 full —
max-error-rows 許容されるエラー行の最大数。エラー行には、 dup-action=STOP_ON_DUP の場合の重複するプライマリキー、カラム数の不一致、および型変換が失敗した行が含まれます。 いいえ 長整数 0 —
timeout バイパスインポートタスク全体のタイムアウト。 いいえ DURATION 7d —
heartbeat-timeout クライアントサイドのハートビートタイムアウト。 いいえ DURATION 60s —
heartbeat-interval クライアントサイドのハートビート間隔。 いいえ DURATION 10s —

型マッピング

MySQL 互換モード

OceanBase 型 Flink 型
TINYINT TINYINT
SMALLINT、TINYINT UNSIGNED SMALLINT
INT、MEDIUMINT、SMALLINT UNSIGNED INT
BIGINT、INT UNSIGNED BIGINT
BIGINT UNSIGNED DECIMAL(20, 0)
REAL、FLOAT FLOAT
DOUBLE DOUBLE
NUMERIC(p, s)、DECIMAL(p, s) DECIMAL(p, s) (p ≤ 38)
BOOLEAN、TINYINT(1) BOOLEAN
DATE DATE
TIME [(p)] TIME [(p)] [WITHOUT TIME ZONE]
DATETIME [(p)]、TIMESTAMP [(p)] TIMESTAMP [(p)] [WITHOUT TIME ZONE]
CHAR(n) CHAR(n)
VARCHAR(n) VARCHAR(n)
BIT(n) BINARY(⌈n/8⌉)
BINARY(n) BINARY(n)
VARBINARY(N) VARBINARY(N)
TINYTEXT、TEXT、MEDIUMTEXT、LONGTEXT STRING
TINYBLOB、BLOB、MEDIUMBLOB、LONGBLOB BYTES (最大 2,147,483,647 バイト)

Oracle 互換モード

OceanBase 型 Flink 型
NUMBER(p, s≤0)、p−s < 3 TINYINT
NUMBER(p, s≤0)、p−s < 5 SMALLINT
NUMBER(p, s≤0)、p−s < 10 INT
NUMBER(p, s≤0)、p−s < 19 BIGINT
NUMBER(p, s≤0)、19 ≤ p−s ≤ 38 DECIMAL(p−s, 0)
NUMBER(p, s>0) DECIMAL(p, s)
NUMBER(p, s≤0)、p−s > 38 STRING
FLOAT、BINARY_FLOAT FLOAT
BINARY_DOUBLE DOUBLE
NUMBER(1) BOOLEAN
DATE、TIMESTAMP [(p)] TIMESTAMP [(p)] [WITHOUT TIME ZONE]
CHAR(n)、NCHAR(n)、NVARCHAR2(n)、VARCHAR(n)、VARCHAR2(n)、CLOB STRING
BLOB、ROWID BYTES (最大 2,147,483,647 バイト)

例

ソーステーブルとシンクテーブル

次の例では、OceanBase のソーステーブルから CDC データを読み取り、JDBC シンクテーブルに書き込みます。参考として、バイパスインポートのシンクテーブル定義も示します。

3 つのテーブルはいずれも 'connector' = 'oceanbase' を使用します。ソースと JDBC シンクでは使用するパラメーターセットが異なります。ダイレクトロードのシンクでは sink.mode = 'direct-load' を設定し、RPC ポート経由で接続します。

-- OceanBase CDC ソーステーブル (全履歴を読み取った後、Binlog を読み取ります)
CREATE TEMPORARY TABLE oceanbase_source (
  a INT,
  b VARCHAR,
  c VARCHAR
) WITH (
  'connector'     = 'oceanbase',
  'hostname'      = '<your-hostname>',
  'port'          = '3306',
  'username'      = '<your-username>',
  'password'      = '<your-password>',
  'database-name' = '<your-database-name>',
  'table-name'    = '<your-table-name>'
);

-- OceanBase JDBC シンクテーブル (リアルタイムのストリーミング書き込み用)
CREATE TEMPORARY TABLE oceanbase_sink (
  a INT,
  b VARCHAR,
  c VARCHAR
) WITH (
  'connector' = 'oceanbase',
  'url'       = '<your-jdbc-url>',
  'userName'  = '<your-username>',
  'password'  = '<your-password>',
  'tableName' = '<your-table-name>'
);

-- OceanBase バイパスインポートのシンクテーブル (高スループットなバッチ書き込み用)
-- 有界データソースが必要です。最適なパフォーマンスを得るには、Flink をバッチモードに設定してください
CREATE TEMPORARY TABLE oceanbase_directload_sink (
  a INT,
  b VARCHAR,
  c VARCHAR
) WITH (
  'connector'   = 'oceanbase',
  'sink.mode'   = 'direct-load',
  'host'        = '<your-host>',
  'port'        = '<your-rpc-port>',
  'tenant-name' = '<your-tenant-name>',
  'schema-name' = '<your-schema-name>',
  'table-name'  = '<your-table-name>',
  'username'    = '<your-username>',
  'password'    = '<your-password>'
);

BEGIN STATEMENT SET;
INSERT INTO oceanbase_sink
SELECT * FROM oceanbase_source;
END;

ディメンションテーブル

次の例では、Datagen ソースを OceanBase のディメンションテーブルとテンポラル結合します。ALL キャッシュポリシーは、ジョブ開始前にディメンションテーブル全体をメモリにロードします。

CREATE TEMPORARY TABLE datagen_source (
  a INT,
  b BIGINT,
  c STRING,
  `proctime` AS PROCTIME()
) WITH (
  'connector' = 'datagen'
);

-- ALL キャッシュポリシーを使用する OceanBase ディメンションテーブル
-- ALL キャッシュは起動時にテーブル全体をロードします。小規模で安定したテーブルに適しています
CREATE TEMPORARY TABLE oceanbase_dim (
  a INT,
  b VARCHAR,
  c VARCHAR
) WITH (
  'connector' = 'oceanbase',
  'url'       = '<your-jdbc-url>',
  'userName'  = '<your-username>',
  'password'  = '${secret_values.password}',
  'tableName' = '<your-table-name>'
);

CREATE TEMPORARY TABLE blackhole_sink (
  a INT,
  b STRING
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN oceanbase_dim FOR SYSTEM_TIME AS OF T.`proctime` AS H
ON T.a = H.a;

次のステップ