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 |
| シンクテーブルでの更新と削除 | はい |
前提条件
開始する前に、次のことを確認してください:
-
ターゲットデータベースとテーブルが OceanBase に存在すること
-
IP アドレスホワイトリストが設定されていること (詳細は「ホワイトリストグループの設定」をご参照ください)
-
(CDC ソーステーブルの場合) OceanBase Binlog サービスが有効になっていること (詳細は「Binlog 関連の操作」をご参照ください)
-
(バイパスインポートシンクテーブルの場合) バイパスインポートポートが有効になっていること (詳細は「バイパスインポート」をご参照ください)
制限
-
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 で指定)。
重要
|
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;
次のステップ
-
サポートされているコネクタ — Realtime Compute for Apache Flink で利用可能なコネクタの完全な一覧
-
OceanBase Binlog service 概要 — CDC 増分読み取りに必要