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

AnalyticDB:Flink CDC による全量・増分データのリアルタイムサブスクライブ (招待制プレビュー)

最終更新日:Aug 25, 2026

AnalyticDB for PostgreSQL は、PostgreSQL の論理レプリケーション機能に基づいて全量データと増分データをサブスクライブする、自社開発の変更データキャプチャ (CDC) コネクタを提供します。このコネクタは Flink とシームレスに統合し、リアルタイムのデータ同期とストリーム処理のために、ソーステーブルからのリアルタイムのデータ変更を効率的にキャプチャします。これにより、企業は動的なデータ要件に迅速に対応できます。このトピックでは、Realtime Compute for Apache Flink CDC を使用して、AnalyticDB for PostgreSQL から全量データと増分データをリアルタイムでサブスクライブする方法について説明します。

制限事項

  • この機能は、マイナーエンジンバージョン 7.2.1.4 以降を実行する AnalyticDB for PostgreSQL V7.0 インスタンスでのみ利用可能です。

    説明

    コンソールのインスタンスの [基本情報] ページでマイナーエンジンバージョンを確認できます。バージョンが上記の要件を満たしていない場合は、マイナーエンジンバージョンをアップグレードする必要があります。

  • AnalyticDB for PostgreSQLのサーバーレスモードはサポートされていません。

前提条件

  • AnalyticDB for PostgreSQL インスタンスとフルマネージド Flink ワークスペースは、同じ VPC 内にある必要があります。

  • AnalyticDB for PostgreSQL インスタンスのパラメーター設定を調整する必要があります:

    • wal_level パラメーターをlogicalに設定して、論理レプリケーションを有効にします。

    • AnalyticDB for PostgreSQL高可用性版インスタンスを使用する場合、hot_standbyhot_standby_feedback、およびsync_replication_slots パラメーターを on に設定する必要があります。これにより、プライマリ/セカンダリフェイルオーバーによって論理サブスクリプションが中断されないようになります。

  • AnalyticDB for PostgreSQL インスタンスには、初期アカウントまたはRDS_SUPERUSER 権限を持つ特権ユーザーを使用する必要があります。 ユーザーには REPLICATION 権限が付与されている必要があります。 ALTER USER <username> WITH REPLICATION;

  • Flink ワークスペースの CIDR ブロックは、AnalyticDB for PostgreSQL インスタンスの ホワイトリスト に追加する必要があります。

  • flink-sql-connector-adbpg-cdc-3.3.jar をダウンロードし、CDC コネクタを Flink ワークスペースにアップロードする必要があります。

操作手順

ステップ 1: テストテーブルとテストデータの準備

  1. AnalyticDB for PostgreSQL コンソールにログインし、対象のインスタンスを見つけます。インスタンス ID をクリックします。

  2. [Basic Information] ページの右下隅にある ログインデータベース をクリックします。

  3. テストデータベースとソーステーブル adbpg_source_table を作成します。次に、ソーステーブルに 50 行のデータを挿入します。

    -- テストデータベースを作成します。
    CREATE DATABASE testdb;
    -- testdb データベースに切り替えてスキーマを作成します。
    \c testdb
    CREATE SCHEMA testschema;
    -- ソーステーブル「adbpg_source_table」を作成します。
    CREATE TABLE testschema.adbpg_source_table(
      id int,
      username text,
      PRIMARY KEY(id)
    );
    -- adbpg_source_table テーブルに 50 行のデータを挿入します。
    INSERT INTO testschema.adbpg_source_table(id, username)
    SELECT i, 'username'||i::text
    FROM generate_series(1, 50) AS t(i);
  4. Flink が結果データを書き込むためのシンクテーブル adbpg_sink_table を作成します。

    CREATE TABLE testschema.adbpg_sink_table(
      id int,
      username text,
      score int
    );

ステップ 2: Flink ジョブの作成

  1. Realtime Compute コンソールにログインします。[Fully Managed Flink] タブで対象のワークスペースを見つけ、[Actions] 列の [Console] をクリックします。

  2. 左側メニューで、[Development] > [ETL] を選択します。

  3. 上部メニューで image をクリックします。[New Blank Stream Draft] を選択し、以下のパラメーターを設定します。

    ジョブパラメーター

    説明

    名前

    ジョブの名前です。

    説明

    ジョブ名は現在のプロジェクト内で一意にする必要があります。

    adbpg-test

    ロケーション

    ジョブのコードファイルが保存されるフォルダーです。

    既存のフォルダーの横の Create a folder アイコンをクリックしてサブフォルダーを作成することもできます。

    Job Drafts

    エンジンバージョン

    ジョブが使用する Flink エンジンのバージョンです。エンジンバージョン番号、バージョンマッピング、ライフサイクルマイルストーンの詳細については、「Flink エンジンバージョン」をご参照ください。

    vvr-6.0.7-flink-1.15

  4. [Create] をクリックします。

ステップ 3: ジョブコードの記述とジョブのデプロイ

  1. アナログデータを生成する datagen_source という名前のソースと、リアルタイムのデータ変更を AnalyticDB for PostgreSQL データベースからキャプチャする source_adbpg という名前のソースを作成します。次に、2 つのソースを結合し、結果を sink_adbpg という名前のシンクテーブルに書き込みます。処理されたデータは AnalyticDB for PostgreSQL に書き込まれます。

    次のジョブコードをエディターにコピーします。

    --- Datagen コネクタを使用してストリーミングデータを生成する Datagen ソーステーブルを作成します。
    CREATE TEMPORARY TABLE datagen_source (
     id INT,
     score INT
    ) WITH (
     'connector' = 'datagen',
     'fields.id.kind'='sequence',
     'fields.id.start'='1',
     'fields.id.end'='100',
     'fields.score.kind'='random',
     'fields.score.min'='70',
     'fields.score.max'='100'
    );
    
    --- adbpg-cdc コネクタを使用して、slot.name と pgoutput に基づいてテーブル「adbpg_source_table」のデータ変更をキャプチャするための adbpg ソーステーブルを作成します。
    CREATE TEMPORARY TABLE source_adbpg(
     id int,
     username varchar,
     PRIMARY KEY(id) NOT ENFORCED
    ) WITH(
      'connector' = 'adbpg-cdc', 
      'hostname' = 'gp-bp16v8cgx46ns****-master.gpdb.rds.aliyuncs.com',
      'port' = '5432',
      'username' = 'account****',
      'password' = 'password****',
      'database-name' = 'testdb',
      'schema-name' = 'testschema',
      'table-name' = 'adbpg_source_table',
      'slot.name' = 'flink',
      'decoding.plugin.name' = 'pgoutput'
    );
    
    --- 処理結果をデータベースの宛先テーブル「adbpg_sink_table」に書き込むための adbpg シンクテーブルを作成します。
    CREATE TEMPORARY TABLE sink_adbpg (
      id int,
      username varchar,
      score int
    ) WITH (
      'connector' = 'adbpg', 
      'url' = 'jdbc:postgresql://gp-bp16v8cgx46ns****-master.gpdb.rds.aliyuncs.com:5432/testdb',
      'tablename' = 'testschema.adbpg_sink_table',  
      'username' = 'account****',
      'password' = 'password****',
      'maxRetryTimes' = '2',
      'batchsize' = '5000',
      'conflictMode' = 'ignore',
      'writeMode' = 'insert',
      'retryWaitTime' = '200'
    );
    
    --- datagen_source テーブルと source_adbpg テーブルの結合結果を adbpg シンクテーブルに書き込みます。
    INSERT INTO sink_adbpg
    SELECT ts.id,ts.username,ds.score FROM datagen_source AS ds
    JOIN source_adbpg AS ts ON ds.id = ts.id;

    パラメーター

    Parameter

    Required

    Data type

    Description

    connector

    Yes

    STRING

    The connector type. Set the value to adbpg-cdc for the source table and adbpg for the sink table.

    hostname

    Yes

    STRING

    The internal endpoint of the AnalyticDB for PostgreSQL instance. You can obtain the internal endpoint on the [Basic Information] page of the instance.

    username

    Yes

    STRING

    The database account and password of the AnalyticDB for PostgreSQL instance.

    password

    Yes

    STRING

    database-name

    Yes

    STRING

    The name of the database.

    schema-name

    Yes

    STRING

    The name of the schema. This parameter supports regular expressions. You can subscribe to multiple schemas at a time.

    table-name

    Yes

    STRING

    The name of the table. This parameter supports regular expressions. You can subscribe to multiple tables at a time.

    port

    Yes

    INTEGER

    The port of AnalyticDB for PostgreSQL. The value is fixed at 5432.

    decoding.plugin.name

    Yes

    STRING

    The name of the PostgreSQL logical decoding plug-in. The value is fixed at pgoutput.

    slot.name

    Yes

    STRING

    The name of the logical decoding slot.

    • For source tables in the same Flink job, use the same value for slot.name.

    • If different Flink jobs involve the same table, set a unique slot.name for each job. This prevents the following error: PSQLException: ERROR: replication slot "debezium" is active for PID 974.

    debezium.*

    No

    STRING

    Controls the behavior of the Debezium client at a finer granularity. For example, setting 'debezium.snapshot.mode' = 'never' disables the snapshot feature. For more information, see the configuration properties.

    scan.incremental.snapshot.enabled

    No

    BOOLEAN

    Specifies whether to enable incremental snapshots. Valid values:

    • false (default): Incremental snapshots are disabled.

    • true: Incremental snapshots are enabled.

    scan.startup.mode

    No

    STRING

    The startup mode for data consumption. Valid values:

    • initial (default): When the job starts for the first time, it scans all historical data and then reads the latest write-ahead logging (WAL) data. This provides a seamless transition between full and incremental data.

    • latest-offset: When the job starts for the first time, it does not scan historical data. It starts reading from the end of the WAL, which is the latest log position. It captures only the data changes that occur after the connector starts.

    • snapshot: Scans all historical data and reads the new WAL entries generated during the full scan. The job stops after the full scan is complete.

    changelog-mode

    No

    STRING

    The changelog mode for encoding stream changes. Valid values:

    • ALL (default): Supports all operation types, including INSERT, DELETE, UPDATE_BEFORE, and UPDATE_AFTER.

    • UPSERT: Supports only UPSERT operations, including INSERT, DELETE, and UPDATE_AFTER.

    heartbeat.interval.ms

    No

    DURATION

    The interval for sending heartbeat packets. The default value is 30 seconds. The unit is milliseconds.

    The AnalyticDB for PostgreSQL CDC connector sends heartbeat packets to the database to ensure that the slot offset continuously advances. If table data does not change frequently, set this parameter to a reasonable value to promptly clean up WAL logs and avoid wasting disk space.

    scan.incremental.snapshot.chunk.key-column

    No

    STRING

    Specifies a column to use for chunking during the snapshot phase. By default, the first column of the primary key is selected.

    url

    Yes

    STRING

    The format is jdbc:postgresql://<Address>:<PortId>/<DatabaseName>.

  2. ジョブ開発ページの上部で [Deep Check] をクリックして構文チェックを実行します。

  3. [Deploy] をクリックし、次に [Confirm] をクリックします。

  4. 右上隅にある [Operations] をクリックします。[Deployments] ページで、[Start] をクリックします。

ステップ 4: Flink が書き込んだデータの表示

  1. テストデータベースで次のステートメントを実行して、Flink が書き込んだデータを表示します。

    SELECT * FROM testschema.adbpg_sink_table;
    SELECT COUNT(*) FROM testschema.adbpg_sink_table; 
  2. ソーステーブルにさらに 50 行のデータを挿入します。次に、Flink がシンクテーブルに書き込む増分データの合計行数を確認します。

    -- ソーステーブルに 50 行の増分データを挿入します。
    INSERT INTO testschema.adbpg_source_table(id, username)
    SELECT i, 'username'||i::text
    FROM generate_series(51, 100) AS t(i);
    
    -- シンクテーブルの新しいデータを確認します。
    SELECT COUNT(*) FROM testschema.adbpg_sink_table where id > 50;

    結果は以下のとおりです。

     count 
    -------
        50
    (1 row)

注意事項

  • ディスク領域の浪費を避けるため、レプリケーションスロットを迅速に管理してください。

    Flink ジョブの再起動中に、チェックポイントに対応する WAL ログがクリーンアップされてデータが失われるのを防ぐため、Flink はレプリケーションスロットを自動的に削除しません。 したがって、Flink ジョブを再起動する必要がなくなったことを確認した場合は、対応するレプリケーションスロットを手動で削除して、占有しているリソースを解放する必要があります。 また、レプリケーションスロットの確認済みの位置が長時間進まない場合、AnalyticDB for PostgreSQL はその位置以降の WAL エントリをクリーンアップできません。 これにより、未使用の WAL データが蓄積され、大量のディスク領域が消費される可能性があります。

  • AnalyticDB for PostgreSQL インスタンスの正常動作時には、厳密に 1 回のデータ処理セマンティクスが保証されます。ただし、障害発生時には、少なくとも 1 回のセマンティクスのみがサポートされます。

  • CDC コネクタは、データ同期の一貫性を確保するために、サブスクライブ対象のテーブルの REPLICA IDENTITY パラメーターを FULL に変更します。この変更には次の影響があります。

    • ディスク領域使用量の増加:更新または削除操作が頻繁に行われるシナリオでは、この設定により WAL ログのサイズが大きくなり、結果としてディスク領域の使用量が増加します。

    • 書き込みパフォーマンスの低下:高並行性の書き込みシナリオでは、パフォーマンスが大幅に影響を受ける可能性があります。

    • チェックポイントの負荷の増加:WAL ログが大きくなると、チェックポイントがより多くのデータを処理する必要があることを意味するため、チェックポイントに必要な時間が長くなる可能性があります。

ベストプラクティス

Flink CDC は、Flink SQL API または DataStream API を使用したジョブ開発をサポートしています。Flink CDC を使用すると、ソースデータベース内の単一または複数のテーブルに対して、フルデータと増分データを統合的に同期できます。また、異種データソースに対するテーブル結合などの計算も実行できます。Flink フレームワークは、データ処理プロセス全体を通じて exactly-once のイベント処理セマンティクスを保証します。ただし、Flink CDC は PostgreSQL 互換データベース全体の同期には適していません。これは、DDL 同期をサポートしておらず、Flink SQL で各テーブルの構造を定義する必要があり、メンテナンスが複雑になるためです。

このセクションでは、AnalyticDB for PostgreSQL から Kafka にデータを同期する例を使用して、Flink CDC SQL ジョブ開発のベストプラクティスについて説明します。Flink CDC ジョブを開発する前に、前提条件のセクションで説明されているように、リソースが準備および設定されていることを確認してください。

ステップ 1: テストテーブルの準備

AnalyticDB for PostgreSQL インスタンスで、2 つのソーステーブルを作成します。

CREATE TABLE products (
    product_id SERIAL PRIMARY KEY,
    product_name VARCHAR(200) NOT NULL,
    sku CHAR(12) NOT NULL,
    description TEXT,
    price NUMERIC(10,2) NOT NULL,
    discount_price DECIMAL(10,2),
    stock_quantity INTEGER DEFAULT 0,
    weight REAL,
    volume DOUBLE PRECISION,
    dimensions BOX,
    release_date DATE,
    is_featured BOOLEAN DEFAULT FALSE,
    rating FLOAT,
    warranty_period INTERVAL,
    metadata JSON,
    tags TEXT[]
);

CREATE TABLE documents (
    document_id UUID PRIMARY KEY,
    title VARCHAR(200) NOT NULL,
    content TEXT,
    summary TEXT,
    publication_date TIMESTAMP WITHOUT TIME ZONE,
    last_updated TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
    author_id BIGINT,
    file_data BYTEA,
    xml_content XML,
    json_metadata JSON,
    reading_time INTERVAL,
    is_public BOOLEAN DEFAULT TRUE,
    views_count INTEGER DEFAULT 0,
    category VARCHAR(50),
    tags TEXT[]
);

ステップ 2: Kafka リソースの準備

  1. Kafka インスタンスを購入してデプロイします

  2. Flink ワークスペースの CIDR ブロックを Kafka インスタンスのホワイトリストに追加します。

  3. Kafka インスタンスでリソースを作成します

ステップ 3: Flink ジョブの作成

  1. Realtime Compute コンソールにログインします。[Fully Managed Flink] タブで対象のワークスペースを見つけ、[Actions] 列の [Console] をクリックします。

  2. 左側メニューで、[Development] > [ETL] を選択します。

  3. 上部メニューで image をクリックします。[New Blank Stream Draft] を選択し、ジョブパラメーターを設定します。

  4. [Create] をクリックします。

ステップ 4: ジョブコードの記述とジョブのデプロイ

  1. Flink ワークスペースで SQL ジョブを記述します。次のジョブコードをエディターにコピーし、設定を実際の値に置き換えます。

    -- 1 つのソースを使用して複数のテーブルからデータをキャプチャします
    CREATE TEMPORARY TABLE ADBPGSource(
        table_name STRING METADATA FROM 'table_name' VIRTUAL,
        row_kind STRING METADATA FROM 'row_kind' VIRTUAL,
        product_id BIGINT,
        product_name STRING,
        sku STRING,
        description STRING,
        price STRING,
        discount_price STRING,
        stock_quantity INT,
        weight STRING,
        volume STRING,
        dimensions STRING,
        release_date STRING,
        is_featured BOOLEAN,
        rating FLOAT,
        warranty_period STRING,
        metadata STRING,
        tags STRING,
        document_id STRING,
        title STRING,
        content STRING,
        summary STRING,
        publication_date STRING,
        last_updated STRING,
        author_id BIGINT,
        file_data STRING,
        xml_content STRING,
        json_metadata STRING,
        reading_time STRING,
        is_public BOOLEAN,
        views_count INT,
        category STRING
    ) WITH (
      'connector' = 'adbpg-cdc',
      'hostname' = 'gp-2zev887z58390***-master.gpdb.rds.aliyuncs.com',
      'port' = '5432',
      'username' = 'account****',
      'password' = 'password****',
      'database-name' = 'testdb',
      'schema-name' = 'public',
      'table-name' = '(products|documents)',
      'slot.name' = 'flink',
      'decoding.plugin.name' = 'pgoutput',
      'debezium.snapshot.mode' = 'never'
    );
    
    CREATE TEMPORARY TABLE KafkaProducts (
        product_id BIGINT,
        product_name STRING,
        sku STRING,
        description STRING,
        price STRING,
        discount_price STRING,
        stock_quantity INT,
        weight STRING,
        volume STRING,
        dimensions STRING,
        release_date STRING,
        is_featured BOOLEAN,
        rating FLOAT,
        warranty_period STRING,
        metadata STRING,
        tags STRING,
        PRIMARY KEY(product_id) NOT ENFORCED
    ) WITH (
      'connector' = 'upsert-kafka',
      'topic' = '****',
      'properties.bootstrap.servers' = 'alikafka-post-cn-****-1-vpc.alikafka.aliyuncs.com:9092',
      'key.format'='avro',
      'value.format'='avro'
    );
    
    CREATE TEMPORARY TABLE KafkaDocuments (
        document_id STRING,
        title STRING,
        content STRING,
        summary STRING,
        publication_date STRING,
        last_updated STRING,
        author_id BIGINT,
        file_data STRING,
        xml_content STRING,
        json_metadata STRING,
        reading_time STRING,
        is_public BOOLEAN,
        views_count INT,
        category STRING,
        tags STRING,
        PRIMARY KEY(document_id) NOT ENFORCED
    ) WITH (
      'connector' = 'upsert-kafka',
      'topic' = '****',
      'properties.bootstrap.servers' = 'alikafka-post-cn-****-1-vpc.alikafka.aliyuncs.com:9092',
      'key.format'='avro',
      'value.format'='avro'
    );
    
    -- STATEMENT SET を使用して複数のステートメントをラップします
    BEGIN STATEMENT SET;
    -- table_name METADATA を使用して、データを宛先テーブルにルーティングします
    INSERT INTO KafkaProducts
    SELECT product_id,product_name,sku,description,price,discount_price,stock_quantity,weight,volume,dimensions,release_date,is_featured,rating,warranty_period,metadata,tags
    FROM ADBPGSource
    WHERE table_name = 'products';
    
    INSERT INTO KafkaDocuments
    SELECT document_id,title,content,summary,publication_date,last_updated,author_id,file_data,xml_content,json_metadata,reading_time,is_public,views_count,category,tags
    FROM ADBPGSource
    WHERE table_name = 'documents';
    
    END;

    この SQL ジョブでは、次の点に注意してください。

    • 複数テーブルの同期タスクでは、この SQL 例のように 1 つのソーステーブルで複数テーブルのデータをキャプチャすることをお勧めします。このソーステーブルには、すべてのソーステーブルのすべての列を定義する必要があります。列名が重複している場合は、1 つだけを残します。宛先テーブルに書き込む際、METADATA table_name を使用して、指定されたテーブルにデータをルーティングします。この方法では、AnalyticDB for PostgreSQL に作成するレプリケーションスロットは 1 つだけで済みます。これにより、ソースデータベースのリソース使用量が削減され、同期パフォーマンスが向上し、将来のメンテナンスが簡素化されます。

    • table-name パラメーターを使用して複数のソーステーブルを指定します。テーブル名を括弧で囲み、縦棒 (|) で区切ります (例:(table1|table2|table3))。

    • debezium.snapshot.modenever に設定すると、ソーステーブルの増分データのみが同期されます。 完全なデータと増分データの両方を同期するには、設定を initial に変更します。

  2. ジョブ開発ページの上部で [Deep Check] をクリックして構文チェックを実行します。

  3. [Deploy] をクリックし、次に [OK] をクリックします。

  4. 右上隅にある [Operations] をクリックします。[Deployments] ページで、[Start] をクリックします。

ステップ 5: テストデータの挿入

AnalyticDB for PostgreSQL インスタンスで、2 つのソーステーブルのデータを更新し、Kafka トピックでメッセージの変更を監視します。

次の SQL ステートメントを使用してテストデータを挿入できます。

INSERT INTO products (
    product_name, sku, description, price, discount_price, stock_quantity, weight, volume, dimensions, release_date, is_featured, rating, warranty_period, metadata, tags
) VALUES (
    'Test Product', 'Test-2025', 'A piece of test product data', 299.99, 279.99, 150, 50.5, 120.75, '(10,20),(30,40)', '2023-05-01', TRUE, 4.8, INTERVAL '1 year', '{"brand": "TechCo", "model": "X1"}', '{"Test1", "Test2"}'
);

関連ドキュメント

AnalyticDB for PostgreSQL のフルデータへのサブスクライブに関する詳細については、「Flink を使用してフルデータをリアルタイムで読み書きする」をご参照ください。