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

PolarDB:PolarDB-O 向け Flink CDC

最終更新日:Aug 27, 2026

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 コネクタをビルドするには、次の手順を実行します。

  1. 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
  2. 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
  3. 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
  4. PolarDB for PostgreSQL (Oracle 互換) との互換性を確保するには、パッチファイルを適用してください。

    git apply release-3.5_support_polardbo.patch
    説明

    この手順で使用するパッチファイルは、こちらからダウンロードできます:release-3.5_support_polardbo.patch

  5. 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 コネクタをビルドするには、次の手順を実行します。

  1. 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
  2. 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
  3. パッチファイルを適用し、PolarDB for PostgreSQL (Oracle 互換) との互換性を確保します。

    git apply release-3.1_support_polardbo.patch
    説明

    この手順で使用するパッチファイルは、こちらからダウンロードできます:release-3.1_support_polardbo.patch

  4. 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 コネクタをビルドするには、次の手順を実行します。

  1. 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
  2. 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
  3. パッチファイルを適用して、PolarDB for PostgreSQL (Oracle 互換)との互換性を確保します。

    git apply release-2.3_support_polardbo.patch
    説明

    この手順で使用するパッチファイルは、こちらからダウンロードできます:release-2.3_support_polardbo.patch

  4. 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 IDENTITYFULL に設定してください。これにより、UPDATE および DELETE イベントにすべての列の古い値が含まれるようになり、データの一貫性を保つために必要です。

    説明
    • REPLICA IDENTITY は PostgreSQL のテーブルレベルの設定であり、論理デコーディングプラグインが UPDATE および DELETE イベントに列の古い値を含めるかどうかを決定します。REPLICA IDENTITY の値の詳細については、「REPLICA IDENTITY」をご参照ください。

    • サブスクライブしているテーブルの REPLICA IDENTITYFULL に設定すると、テーブルロックが必要になり、業務に影響を与える可能性があります。この点を考慮して計画してください。次のコマンドを使用して、現在の設定が 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

  • PolarDB for PostgreSQL (Oracle 互換) 1.0: TIMESTAMP

  • PolarDB for PostgreSQL (Oracle 互換) 2.0: TIMESTAMP

  • PolarDB for PostgreSQL: 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 ドキュメントを参照してコネクタのパラメーターを設定してください。

  1. 前提条件

    • PolarDB for PostgreSQL (Compatible with Oracle) の準備

      1. PolarDB クラスターの購入ページで、PolarDB for PostgreSQL (Compatible with Oracle) 2.0 クラスターを購入します。

      2. 特権アカウントを作成します

      3. クラスターのプライマリエンドポイントを表示しますPolarDB クラスターと Realtime Compute for Apache Flink ワークスペースが同じ仮想プライベートクラウド (VPC) 内にある場合は、プライベートエンドポイントを使用できます。それ以外の場合は、パブリックエンドポイントを申請して使用する必要があります。

      4. クラスターの IP アドレスホワイトリストを設定します。Flink インスタンスの IP アドレスを PolarDB クラスターのホワイトリストに追加します。

      5. コンソールで、ソースデータベース flink_source と宛先データベース flink_sink を作成します。詳細については、「データベースの作成」をご参照ください。

      6. 次のステートメントを実行して、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();
      7. 次のステートメントを実行して、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 の準備

      1. Realtime Compute コンソールにログオンし、Realtime Compute for Apache Flink インスタンスを購入します。詳細については、「Realtime Compute for Apache Flink のアクティベート」をご参照ください。

        説明

        Realtime Compute for Apache Flink ワークスペースは、PolarDB クラスターと同じ [リージョン] および [VPC] に作成することを推奨します。これにより、接続に PolarDB クラスターのプライベートプライマリエンドポイントを使用できます。

      2. カスタムコネクタを作成し、ビルドした PolarDB-O Flink CDC パッケージをアップロードします。[Formats] には debezium-json を選択します。詳細については、「カスタムコネクタの作成」をご参照ください。

  2. Flink ジョブの作成

    1. 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;
    2. ジョブのデプロイと開始。

      ジョブエディターの上部ツールバーで、[Deploy] をクリックします。

      左側メニューで、[Operations Center] > [Deployments] を選択します。ジョブリストでターゲットジョブを見つけ、[操作] 列の [Start] をクリックします。

    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       | 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 パイプラインコネクタ のドキュメントを参照してパラメーターを設定してください。

  1. 前提条件

    • PolarDB for PostgreSQL (Compatible with Oracle) の準備

      1. PolarDB クラスターの購入ページで、PolarDB for PostgreSQL (Compatible with Oracle) 2.0 クラスターを購入します。

      2. 特権アカウントを作成します

      3. クラスターのプライマリエンドポイントを表示しますPolarDB クラスターと Realtime Compute for Apache Flink ワークスペースが同じ VPC 内にある場合は、プライベートエンドポイントを使用できます。それ以外の場合は、パブリックエンドポイントを申請して使用する必要があります。

      4. クラスターの IP アドレスホワイトリストを設定します。Flink インスタンスの IP アドレスを PolarDB クラスターのホワイトリストに追加します。

      5. コンソールで、ソースデータベース flink_source を作成します。詳細については、「データベースの作成」をご参照ください。

      6. 次のステートメントを実行して、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 クラスターのプライベートプライマリエンドポイントを使用できます。

  2. Flink ジョブの作成

    1. 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
    2. 左側の More セクションに、作成したパイプラインコネクタを追加します。右側の [詳細設定] パネルで、エンジンバージョンが vvr-11.5-jdk11-flink-1.20 であることを確認し、[追加の依存ファイル] 領域に必要な依存ファイルを追加します。

    3. ジョブのデプロイと開始。

      1. 右上隅のデプロイメントをクリックします。

        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
      2. 「デプロイメント」ページに移動し、Enable をクリックします。

    4. 結果のテストと検証。

      • デプロイメントジョブが正常に実行されると、ステータスは「実行中」になります。全データフェーズの 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=()}