PolarDB for PostgreSQL (Compatible with Oracle) 向けの Flink CDC コネクタ (PolarDB-O Flink CDC コネクタ) は、PolarDB for PostgreSQL (Compatible with Oracle) データベースから完全なデータスナップショットとその後の増分変更を読み取ります。機能と使用方法については、コミュニティの Postgres CDC ドキュメントをご参照ください。
PolarDB for PostgreSQL (Compatible with Oracle) とコミュニティ PostgreSQL は、データ型と組み込みオブジェクトの処理においてわずかな違いしかないため、この記事では、わずかなコード変更でコミュニティ Postgres CDC コネクターを適合させ、PolarDB for PostgreSQL (Compatible with Oracle) 用の PolarDB Flink CDC コネクターをパッケージ化する方法について説明します。
PolarDB for PostgreSQL (Oracle 互換)の DATE 型は 64 ビットですが、コミュニティ PostgreSQL では 32 ビットです。PolarDB-O Flink CDC コネクタがこの違いを処理します。
PolarDB-O Flink CDC のビルド
PolarDB-O Flink CDC コネクタは、コミュニティの Postgres CDC をベースにしています。このコネクタには、ご自身でビルドする場合でも、このトピックで提供されている JAR パッケージを使用する場合でも、サービスレベルアグリーメント (SLA) は提供されません。
前提条件
-
Flink-CDC バージョンの決定
Alibaba Cloud Realtime Compute for Apache Flink を使用する場合、お使いの Ververica Runtime (VVR) バージョンに対応するコミュニティ Flink-CDC バージョンを特定する必要があります。詳細については、「CDC と VVR のバージョンマッピング」をご参照ください。
説明Flink-CDC のコードリポジトリは Flink-CDC にあります。
-
Debezium バージョンの決定
対応する Flink-CDC バージョンの
pom.xmlファイルで、debezium.versionプロパティを見つけて、Debezium バージョンを決定します。説明Debezium のコードリポジトリは Debezium にあります。
-
PgJDBC バージョンの決定
対応する Postgres-CDC バージョンの
pom.xmlファイルで、org.postgresql依存関係を見つけて、PgJDBC バージョンを決定します。説明-
release-3.0 未満のバージョンの場合、ファイルパスは
flink-connector-postgres-cdc/pom.xmlです。 -
release-3.0 以降の場合、ファイルパスは
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/pom.xmlです。 -
PgJDBC のコードリポジトリは PgJDBC にあります。
-
操作手順
release-3.5 のビルド
Community Flink-CDC release-3.5 は、Alibaba Cloud Realtime Compute for Apache Flink の vvr-11.4-jdk11-flink-1.20 と互換性があります。
このバージョンの PolarDB-O Flink CDC コネクタをビルドするには、次の手順を実行します。
-
Flink-CDC、Debezium、PgJDBC の対応するバージョンのリポジトリをクローンします。
git clone -b release-3.5 --depth=1 https://github.com/apache/flink-cdc.git git clone -b REL42.7.3 --depth=1 https://github.com/pgjdbc/pgjdbc.git git clone -b v1.9.8.Final --depth=1 https://github.com/debezium/debezium.git -
Debezium および PgJDBC リポジトリから必要なファイルを Flink-CDC ディレクトリにコピーします。
mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/Oid.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core cp debezium/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql -
Flink-CDC ディレクトリに移動し、タイムスタンプ変換のバグ修正と暗黙的な動作 (SELECT *) のバグ修正 を適用します。これらの修正は、コミュニティの release 3.6 にマージされる予定です。
cd flink-cdc # タイムスタンプ変換のバグ修正を適用します。 git fetch origin 2f32836a783f80f295c9dce339c11afec2a32dc2 git cherry-pick 2f32836a783f80f295c9dce339c11afec2a32dc2 # 暗黙的な動作 (SELECT *) のバグ修正を適用します。 git fetch origin 0d86de24494a855c2d83f9b1052c2e888e182cb1 git cherry-pick 0d86de24494a855c2d83f9b1052c2e888e182cb1 -
PolarDB for PostgreSQL (Oracle 互換) との互換性を確保するには、パッチファイルを適用してください。
git apply release-3.5_support_polardbo.patch説明この手順で使用するパッチファイルは、こちらからダウンロードできます:release-3.5_support_polardbo.patch
-
Maven を使用して PolarDB-O Flink CDC コネクタをビルドします。
mvn clean install -DskipTests -Dcheckstyle.skip=true -Dspotless.check.skip # ビルドが完了すると、JAR パッケージは flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-postgres/target ディレクトリに生成されます。
上記の手順に従って JDK 11 でビルドされた PolarDB-O Flink CDC コネクタの JAR パッケージは次のとおりです:flink-cdc-pipeline-connector-polardbo-3.5-SNAPSHOT-20260212.jar
release-3.1 のビルド
コミュニティ Flink-CDC リリース 3.1 は、Alibaba Cloud Realtime Compute for Apache Flink の vvr-8.0.x-flink-1.17 と互換性があります。
このバージョンの PolarDB-O Flink CDC コネクタをビルドするには、次の手順を実行します。
-
Flink-CDC、Debezium、PgJDBC の対応するバージョンのリポジトリをクローンします。
git clone -b release-3.1 --depth=1 https://github.com/apache/flink-cdc.git git clone -b REL42.5.1 --depth=1 https://github.com/pgjdbc/pgjdbc.git git clone -b v1.9.8.Final --depth=1 https://github.com/debezium/debezium.git -
Debezium および PgJDBC リポジトリから必要なファイルを Flink-CDC ディレクトリにコピーします。
mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp debezium/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql -
パッチファイルを適用し、PolarDB for PostgreSQL (Oracle 互換) との互換性を確保します。
git apply release-3.1_support_polardbo.patch説明この手順で使用するパッチファイルは、こちらからダウンロードできます:release-3.1_support_polardbo.patch
-
Maven を使用して PolarDB-O Flink CDC コネクタをビルドします。
mvn clean install -DskipTests -Dcheckstyle.skip=true -Dspotless.check.skip -Drat.skip=true # ビルドが完了すると、JAR パッケージは flink-sql-connector-postgres-cdc/target ディレクトリに生成されます。
上記の手順に従って JDK 8 でビルドされた PolarDB-O Flink CDC コネクタの JAR パッケージは次のとおりです:flink-sql-connector-postgres-cdc-3.1-SNAPSHOT.jar
release-2.3 のビルド
コミュニティ Flink-CDC release-2.3 は、Alibaba Cloud Realtime Compute for Apache Flink の vvr-4.0.15-flink-1.13 から vvr-6.0.2-flink-1.15 までと互換性があります。
このバージョンの PolarDB-O Flink CDC コネクタをビルドするには、次の手順を実行します。
-
Flink-CDC、Debezium、PgJDBC の対応するバージョンのリポジトリをクローンします。
git clone -b release-2.3 --depth=1 https://github.com/apache/flink-cdc.git git clone -b REL42.2.26 --depth=1 https://github.com/pgjdbc/pgjdbc.git git clone -b v1.6.4.Final --depth=1 https://github.com/debezium/debezium.git -
Debezium および PgJDBC リポジトリから必要なファイルを Flink-CDC ディレクトリにコピーします。
mkdir -p flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 mkdir -p flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc mkdir -p flink-cdc/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp debezium/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java flink-cdc/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql -
パッチファイルを適用して、PolarDB for PostgreSQL (Oracle 互換)との互換性を確保します。
git apply release-2.3_support_polardbo.patch説明この手順で使用するパッチファイルは、こちらからダウンロードできます:release-2.3_support_polardbo.patch
-
Maven を使用して PolarDB-O Flink CDC コネクタをビルドします。
mvn clean install -DskipTests -Dcheckstyle.skip=true -Dspotless.check.skip -Drat.skip=true # ビルドが完了すると、JAR パッケージは flink-sql-connector-postgres-cdc/target ディレクトリに生成されます。
上記の手順に従って JDK 8 でビルドされた PolarDB-O Flink CDC コネクタの JAR パッケージは次のとおりです:flink-sql-connector-postgres-cdc-2.3-SNAPSHOT.jar
使用方法
PolarDB-O Flink CDC コネクタは、論理レプリケーションを使用して PolarDB for PostgreSQL (Compatible with Oracle) データベースから CDC データストリームを読み取ります。これには、以下の要件があります。
-
wal_levelパラメーターをlogicalに設定してください。この設定により、論理レプリケーションに必要な情報が先行書き込みログ (WAL) ファイルに追加されます。説明wal_level パラメーターはコンソールで設定できます。詳細な手順については、「クラスターパラメーターの設定」をご参照ください。このパラメーターを変更すると、クラスターが再起動します。業務への影響を最小限に抑えるために、この操作は慎重に計画してください。
-
ALTER TABLE schema.table REPLICA IDENTITY FULL;コマンドを実行して、サブスクライブしているテーブルのREPLICA IDENTITYをFULLに設定してください。これにより、UPDATE および DELETE イベントにすべての列の古い値が含まれるようになり、データの一貫性を保つために必要です。説明-
REPLICA IDENTITY は PostgreSQL のテーブルレベルの設定であり、論理デコーディングプラグインが UPDATE および DELETE イベントに列の古い値を含めるかどうかを決定します。REPLICA IDENTITY の値の詳細については、「REPLICA IDENTITY」をご参照ください。
-
サブスクライブしているテーブルの
REPLICA IDENTITYをFULLに設定すると、テーブルロックが必要になり、業務に影響を与える可能性があります。この点を考慮して計画してください。次のコマンドを使用して、現在の設定がFULLであるかどうかを確認できます。SELECT relreplident = 'f' FROM pg_class WHERE relname = 'tablename';
-
-
max_wal_senders および max_replication_slots パラメーターの値が、現在使用されているレプリケーションスロットの数と Flink ジョブで必要なスロット数の合計を超えることを確認してください。
-
特権アカウント、または LOGIN および REPLICATION の両方の権限を持つアカウントを使用してください。このアカウントには、初期スナップショットクエリを実行するために、サブスクライブしているテーブルに対する SELECT 権限も必要です。
-
PolarDB クラスターのプライマリエンドポイントにのみ接続できます。クラスターエンドポイントは論理レプリケーションをサポートしていません。
-
Release 3.5 以降では、親テーブルを指定してパーティションテーブルを同期できます。以下の設定が必要です。詳細については、「Postgres CDC コミュニティのドキュメント」をご参照ください。
-
scan.include-partitioned-tables.enabledオプションをtrueに設定してください。 -
publish_via_partition_root=trueオプションを使用して、データベースにPUBLICATIONを手動で作成します。次に、debezium.publication.nameパラメーターを使用してtable-nameを指定します。 -
Flink コネクタの
table-nameオプションには親テーブルのみを指定する必要があります。スナップショットフェーズでデータが重複するのを避けるため、正規表現が子テーブルと一致しないようにしてください。
さらに、Release 3.5 以降では、スナップショットデータと増分データを読み取り、エンドツーエンドの完全なデータベース同期を提供するパイプラインコネクタもサポートしています。ただし、パイプラインコネクタは現在スキーマの変更をサポートしていません。詳細については、「Postgres CDC パイプラインコネクタのコミュニティドキュメント」をご参照ください。
-
PolarDB-O Flink CDC と Postgres CDC の違い
PolarDB-O Flink CDC コネクタは、Postgres CDC コネクタに基づいて構築されています。構文とパラメーターについては、Postgres CDC のドキュメントをご参照ください。ただし、主な違いは次のとおりです。
-
WITH 句では、'connector' パラメーターを固定値
polardbo-cdcに設定する必要があります。 -
PolarDB-O Flink CDC は、PolarDB for PostgreSQL のすべてのバージョン、PolarDB for PostgreSQL (Compatible with Oracle) 1.0、および PolarDB for PostgreSQL (Compatible with Oracle) 2.0 と互換性があります。
説明PolarDB for PostgreSQL を使用している場合は、コミュニティの Postgres CDC コネクタを直接使用することを推奨します。
-
PolarDB for PostgreSQL (Compatible with Oracle) 1.0 および PolarDB for PostgreSQL (Compatible with Oracle) 2.0 の
DATE型の列については、Flink SQL のソーステーブルとシンクテーブルの両方で、対応する型をTIMESTAMPとして指定する必要があります。 -
decoding.plugin.nameパラメーターはpgoutputに設定することを推奨します。設定しない場合、UTF-8 以外のエンコーディングを使用するデータベースでは、増分解析中に文字化けが発生する可能性があります。詳細については、「コミュニティドキュメント」をご参照ください。
データ型マッピング
PolarDB for PostgreSQL と Flink 間のデータ型のマッピングは、DATE 型を除き、コミュニティ PostgreSQL のものと同一です。 完全なマッピングは次のとおりです。
|
ソース タイプ |
Flink タイプ |
|
SMALLINT |
SMALLINT |
|
INT2 |
|
|
SMALLSERIAL |
|
|
SERIAL2 |
|
|
INTEGER |
INT |
|
SERIAL |
|
|
BIGINT |
BIGINT |
|
BIGSERIAL |
|
|
REAL |
FLOAT |
|
FLOAT4 |
|
|
FLOAT8 |
DOUBLE |
|
DOUBLE PRECISION |
|
|
NUMERIC(p, s) |
DECIMAL(p, s) |
|
DECIMAL(p, s) |
|
|
BOOLEAN |
BOOLEAN |
|
DATE |
|
|
TIME [(p)] [WITHOUT TIMEZONE] |
TIME [(p)] [WITHOUT TIMEZONE] |
|
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
|
CHAR(n) |
STRING |
|
CHARACTER(n) |
|
|
VARCHAR(n) |
|
|
CHARACTER VARYING(n) |
|
|
TEXT |
|
|
BYTEA |
BYTES |
例
ソースコネクタ
この例では、PolarDB-O Flink CDC コネクタを使用して、PolarDB for PostgreSQL (Compatible with Oracle) 2.0 クラスター上の flink_source データベースから shipments テーブルを flink_sink データベースの shipments_sink テーブルに同期する方法を示します。
この例は、PolarDB for PostgreSQL (Compatible with Oracle) で PolarDB-O Flink CDC コネクタを実行する基本的なデモです。本番環境で使用する場合は、ビジネス要件に基づいて、コミュニティの Postgres CDC ドキュメントを参照してコネクタのパラメーターを設定してください。
-
前提条件
-
PolarDB for PostgreSQL (Compatible with Oracle) の準備
-
PolarDB クラスターの購入ページで、PolarDB for PostgreSQL (Compatible with Oracle) 2.0 クラスターを購入します。
-
クラスターのプライマリエンドポイントを表示します。PolarDB クラスターと Realtime Compute for Apache Flink ワークスペースが同じ仮想プライベートクラウド (VPC) 内にある場合は、プライベートエンドポイントを使用できます。それ以外の場合は、パブリックエンドポイントを申請して使用する必要があります。
-
クラスターの IP アドレスホワイトリストを設定します。Flink インスタンスの IP アドレスを PolarDB クラスターのホワイトリストに追加します。
-
コンソールで、ソースデータベース flink_source と宛先データベース flink_sink を作成します。詳細については、「データベースの作成」をご参照ください。
-
次のステートメントを実行して、flink_source データベースに shipments テーブルを作成し、データを挿入します。
CREATE TABLE public.shipments ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time DATE, PRIMARY KEY (shipment_id) ); ALTER TABLE public.shipments REPLICA IDENTITY FULL; INSERT INTO public.shipments SELECT 1, 1, 'test1', 'test1', false, now(); -
次のステートメントを実行して、flink_sink データベースに shipments_sink テーブルを作成します。
CREATE TABLE public.shipments_sink ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time TIMESTAMP, PRIMARY KEY (shipment_id) );
-
-
Realtime Compute for Apache Flink の準備
-
Realtime Compute コンソールにログオンし、Realtime Compute for Apache Flink インスタンスを購入します。詳細については、「Realtime Compute for Apache Flink のアクティベート」をご参照ください。
説明Realtime Compute for Apache Flink ワークスペースは、PolarDB クラスターと同じ [リージョン] および [VPC] に作成することを推奨します。これにより、接続に PolarDB クラスターのプライベートプライマリエンドポイントを使用できます。
-
カスタムコネクタを作成し、ビルドした PolarDB-O Flink CDC パッケージをアップロードします。[Formats] には debezium-json を選択します。詳細については、「カスタムコネクタの作成」をご参照ください。
-
-
-
Flink ジョブの作成
-
Realtime Compute for Apache Flink コンソールにログオンし、SQL ドラフトを作成します。詳細については、「SQL ドラフトの開発」をご参照ください。次の Flink SQL コードを使用し、PolarDB クラスターのプライマリエンドポイント、ポート、ユーザー名、パスワードのプレースホルダーを置き換えてください。
説明PolarDB for PostgreSQL (Oracle 互換) の DATE 型は 64 ビットですが、Flink SQL およびほとんどのデータベースの DATE 型は 32 ビットです。したがって、ソーステーブルの
DATE列を、Flink SQL のソーステーブルとシンクテーブルの両方で TIMESTAMP 型にマッピングする必要があります。そうしないと、ジョブは"java.time.DateTimeException: Invalid value for EpochDay (valid values -365243219162 - 365241780471): 1720891573000"のような型の不一致エラーで失敗します。CREATE TEMPORARY TABLE shipments ( shipment_id INT, order_id INT, origin STRING, destination STRING, is_arrived BOOLEAN, order_time TIMESTAMP, PRIMARY KEY (shipment_id) NOT ENFORCED ) WITH ( 'connector' = 'polardbo-cdc', 'hostname' = '<yourHostname>', 'port' = '<yourPort>', 'username' = '<yourUserName>', 'password' = '<yourPassWord>', 'database-name' = 'flink_source', 'schema-name' = 'public', 'table-name' = 'shipments', 'decoding.plugin.name' = 'pgoutput', 'slot.name' = 'flink' ); CREATE TEMPORARY TABLE shipments_sink ( shipment_id INT, order_id INT, origin STRING, destination STRING, is_arrived BOOLEAN, order_time TIMESTAMP, PRIMARY KEY (shipment_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:postgresql://<yourHostname>:<yourPort>/flink_sink', 'table-name' = 'shipments_sink', 'username' = '<yourUserName>', 'password' = '<yourPassWord>' ); INSERT INTO shipments_sink SELECT * FROM shipments; -
ジョブのデプロイと開始。
ジョブエディターの上部ツールバーで、[Deploy] をクリックします。
左側メニューで、[Operations Center] > [Deployments] を選択します。ジョブリストでターゲットジョブを見つけ、[操作] 列の [Start] をクリックします。
-
結果のテストと検証。
-
ジョブがデプロイされ、実行中の状態になると、shipments テーブルのデータが flink_sink データベースの shipments_sink テーブルに同期されます。
SELECT * FROM public.shipments_sink;次の結果が返されます:
shipment_id | order_id | origin | destination | is_arrived | order_time -------------+----------+--------+-------------+------------+--------------------- 1 | 1 | test1 | test1 | f | 2024-09-18 05:45:08 (1 row) -
flink_source データベースの shipments テーブルで DML ステートメントを実行します。変更はリアルタイムで同期されます。
INSERT INTO public.shipments SELECT 2, 2, 'test2', 'test2', false, now(); UPDATE public.shipments SET is_arrived = true WHERE shipment_id = 1; DELETE FROM public.shipments WHERE shipment_id = 2; INSERT INTO public.shipments SELECT 3, 3, 'test3', 'test3', false, now(); UPDATE public.shipments SET is_arrived = true WHERE shipment_id = 3;shipments テーブルのデータが flink_sink データベースの shipments_sink テーブルに同期されます。
SELECT * FROM public.shipments_sink;次の結果が返されます:
shipment_id | order_id | origin | destination | is_arrived | order_time -------------+----------+--------+-------------+------------+--------------------- 1 | 1 | test1 | test1 | t | 2024-09-18 05:45:08 3 | 3 | test3 | test3 | t | 2024-09-18 07:33:23 (2 rows)
-
-
パイプラインコネクタ
この例では、PolarDB-O Flink CDC パイプラインコネクタを使用して、PolarDB for PostgreSQL (Compatible with Oracle) 2.0 クラスターから shipments1 および shipments2 テーブルを同期する方法を示します。デバッグのため、シンクには Print コネクタ を使用します。本番環境では、ビジネス要件に基づいて適切なシンクコネクタを選択してください。
この例は、PolarDB for PostgreSQL (Compatible with Oracle) で PolarDB-O Flink CDC コネクタを実行する基本的なデモです。本番環境で使用する場合は、ビジネス要件に基づいて、コミュニティの Postgres CDC パイプラインコネクタ のドキュメントを参照してパラメーターを設定してください。
-
前提条件
-
PolarDB for PostgreSQL (Compatible with Oracle) の準備
-
PolarDB クラスターの購入ページで、PolarDB for PostgreSQL (Compatible with Oracle) 2.0 クラスターを購入します。
-
クラスターのプライマリエンドポイントを表示します。PolarDB クラスターと Realtime Compute for Apache Flink ワークスペースが同じ VPC 内にある場合は、プライベートエンドポイントを使用できます。それ以外の場合は、パブリックエンドポイントを申請して使用する必要があります。
-
クラスターの IP アドレスホワイトリストを設定します。Flink インスタンスの IP アドレスを PolarDB クラスターのホワイトリストに追加します。
-
コンソールで、ソースデータベース flink_source を作成します。詳細については、「データベースの作成」をご参照ください。
-
次のステートメントを実行して、flink_source データベースに shipments1 および shipments2 テーブルを作成し、データを挿入します。
CREATE TABLE public.shipments1 ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time DATE, PRIMARY KEY (shipment_id) ); ALTER TABLE public.shipments1 REPLICA IDENTITY FULL; INSERT INTO public.shipments1 SELECT 1, 1, 'test1', 'test1', false, now(); CREATE TABLE public.shipments2 ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time DATE, PRIMARY KEY (shipment_id) ); ALTER TABLE public.shipments2 REPLICA IDENTITY FULL; INSERT INTO public.shipments2 SELECT 1, 1, 'test1', 'test1', false, now();
-
-
Realtime Compute for Apache Flink の準備
Realtime Compute コンソールにログオンし、Realtime Compute for Apache Flink インスタンスを購入します。詳細については、「Realtime Compute for Apache Flink のアクティベート」をご参照ください。
説明Realtime Compute for Apache Flink ワークスペースは、PolarDB クラスターと同じ [リージョン] および [VPC] に作成することを推奨します。これにより、接続に PolarDB クラスターのプライベートプライマリエンドポイントを使用できます。
-
-
Flink ジョブの作成
-
Realtime Compute for Apache Flink コンソールにログオンし、データ取り込みドラフトを作成します。詳細については、「Flink CDC データ取り込み」をご参照ください。次のデータ取り込み設定を使用し、PolarDB クラスターのプライマリエンドポイント、ポート、ユーザー名、パスワードのプレースホルダーを置き換えてください。
source: type: polardbo name: PolarDB Oracle Source hostname: '<yourHostname>' port: '<yourPort>' username: '<yourUserName>' password: '<yourPassWord>' tables: flink_source.public.shipments[12] decoding.plugin.name: pgoutput slot.name: pgtest sink: type: values name: values Sink print.enabled: true -
左側の More セクションに、作成したパイプラインコネクタを追加します。右側の [詳細設定] パネルで、エンジンバージョンが
vvr-11.5-jdk11-flink-1.20であることを確認し、[追加の依存ファイル] 領域に必要な依存ファイルを追加します。 -
ジョブのデプロイと開始。
-
右上隅のデプロイメントをクリックします。
source: type: polardbo name: PolarDB Oracle Source hostname: xxx port: xxx username: xxx password: xxx tables: flink_source.public.shipments[12] decoding.plugin.name: pgoutput slot.name: pgtest sink: type: values name: values Sink print.enabled: true -
「デプロイメント」ページに移動し、Enable をクリックします。
-
-
結果のテストと検証。
-
デプロイメントジョブが正常に実行されると、ステータスは「実行中」になります。全データフェーズの CreateTableEvent と DataChangeEvent は、[ジョブログ] > [実行中のタスクマネージャー] > [Stdout] ログで確認できます。デプロイメント詳細ページで、[ジョブログ] タブをクリックし、[実行中のタスクマネージャー] を選択し、次に [Stdout] タブをクリックしてログ出力を表示します。出力に次のテーブル作成イベントとデータ変更イベントが含まれていることを確認します。
CreateTableEvent{tableId=public.shipments2, schema=columns={`shipment_id` INT NOT NULL,`order_id` INT,`origin` STRING,`destination` STRING,`is_arrived` BOOLEAN,`order_time` TIMESTAMP(6)}, primaryKeys=shipment_id, options=()} CreateTableEvent{tableId=public.shipments1, schema=columns={`shipment_id` INT NOT NULL,`order_id` INT,`origin` STRING,`destination` STRING,`is_arrived` BOOLEAN,`order_time` TIMESTAMP(6)}, primaryKeys=shipment_id, options=()} DataChangeEvent{tableId=public.shipments2, before=[], after=[1, 1, test1, test1, false, 2026-01-07T16:30:44], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments1, before=[], after=[1, 1, test1, test1, false, 2026-01-07T16:30:44], op=INSERT, meta=()} -
flink_source データベースの shipments1 および shipments2 テーブルで DML ステートメントを実行します。変更はリアルタイムで同期されます。
INSERT INTO public.shipments1 SELECT 2, 2, 'test2', 'test2', false, now(); UPDATE public.shipments1 SET is_arrived = true WHERE shipment_id = 1; DELETE FROM public.shipments1 WHERE shipment_id = 2; INSERT INTO public.shipments1 SELECT 3, 3, 'test3', 'test3', false, now(); UPDATE public.shipments1 SET is_arrived = true WHERE shipment_id = 3; INSERT INTO public.shipments2 SELECT 2, 2, 'test2', 'test2', false, now(); UPDATE public.shipments2 SET is_arrived = true WHERE shipment_id = 1; DELETE FROM public.shipments2 WHERE shipment_id = 2; INSERT INTO public.shipments2 SELECT 3, 3, 'test3', 'test3', false, now(); UPDATE public.shipments2 SET is_arrived = true WHERE shipment_id = 3; -
増分フェーズの DataChangeEvent は、[ジョブログ] > [実行中のタスクマネージャー] > [Stdout] ログで確認できます:
DataChangeEvent{tableId=public.shipments1, before=[], after=[2, 2, test2, test2, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments1, before=[1, 1, test1, test1, false, 2026-01-07T16:30:44], after=[1, 1, test1, test1, true, 2026-01-07T16:30:44], op=UPDATE, meta=()} DataChangeEvent{tableId=public.shipments1, before=[2, 2, test2, test2, false, 2026-01-07T16:44:50], after=[], op=DELETE, meta=()} DataChangeEvent{tableId=public.shipments1, before=[], after=[3, 3, test3, test3, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments1, before=[3, 3, test3, test3, false, 2026-01-07T16:44:50], after=[3, 3, test3, test3, true, 2026-01-07T16:44:50], op=UPDATE, meta=()} DataChangeEvent{tableId=public.shipments2, before=[], after=[2, 2, test2, test2, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments2, before=[1, 1, test1, test1, false, 2026-01-07T16:30:44], after=[1, 1, test1, test1, true, 2026-01-07T16:30:44], op=UPDATE, meta=()} DataChangeEvent{tableId=public.shipments2, before=[2, 2, test2, test2, false, 2026-01-07T16:44:50], after=[], op=DELETE, meta=()} DataChangeEvent{tableId=public.shipments2, before=[], after=[3, 3, test3, test3, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments2, before=[3, 3, test3, test3, false, 2026-01-07T16:44:50], after=[3, 3, test3, test3, true, 2026-01-07T16:44:50], op=UPDATE, meta=()}
-
-