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_standby、hot_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: テストテーブルとテストデータの準備
-
AnalyticDB for PostgreSQL コンソールにログインし、対象のインスタンスを見つけます。インスタンス ID をクリックします。
-
[Basic Information] ページの右下隅にある ログインデータベース をクリックします。
-
テストデータベースとソーステーブル 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); -
Flink が結果データを書き込むためのシンクテーブル adbpg_sink_table を作成します。
CREATE TABLE testschema.adbpg_sink_table( id int, username text, score int );
ステップ 2: Flink ジョブの作成
-
Realtime Compute コンソールにログインします。[Fully Managed Flink] タブで対象のワークスペースを見つけ、[Actions] 列の [Console] をクリックします。
-
左側メニューで、 を選択します。
-
上部メニューで
をクリックします。[New Blank Stream Draft] を選択し、以下のパラメーターを設定します。ジョブパラメーター
説明
例
名前
ジョブの名前です。
説明ジョブ名は現在のプロジェクト内で一意にする必要があります。
adbpg-test
ロケーション
ジョブのコードファイルが保存されるフォルダーです。
既存のフォルダーの横の
アイコンをクリックしてサブフォルダーを作成することもできます。Job Drafts
エンジンバージョン
ジョブが使用する Flink エンジンのバージョンです。エンジンバージョン番号、バージョンマッピング、ライフサイクルマイルストーンの詳細については、「Flink エンジンバージョン」をご参照ください。
vvr-6.0.7-flink-1.15
-
[Create] をクリックします。
ステップ 3: ジョブコードの記述とジョブのデプロイ
-
アナログデータを生成する
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-cdcfor the source table andadbpgfor 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.namefor 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, andUPDATE_AFTER. -
UPSERT: Supports only
UPSERToperations, includingINSERT,DELETE, andUPDATE_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>. -
-
ジョブ開発ページの上部で [Deep Check] をクリックして構文チェックを実行します。
-
[Deploy] をクリックし、次に [Confirm] をクリックします。
-
右上隅にある [Operations] をクリックします。[Deployments] ページで、[Start] をクリックします。
ステップ 4: Flink が書き込んだデータの表示
-
テストデータベースで次のステートメントを実行して、Flink が書き込んだデータを表示します。
SELECT * FROM testschema.adbpg_sink_table; SELECT COUNT(*) FROM testschema.adbpg_sink_table; -
ソーステーブルにさらに 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 リソースの準備
-
Flink ワークスペースの CIDR ブロックを Kafka インスタンスのホワイトリストに追加します。
ステップ 3: Flink ジョブの作成
-
Realtime Compute コンソールにログインします。[Fully Managed Flink] タブで対象のワークスペースを見つけ、[Actions] 列の [Console] をクリックします。
-
左側メニューで、 を選択します。
-
上部メニューで
をクリックします。[New Blank Stream Draft] を選択し、ジョブパラメーターを設定します。 -
[Create] をクリックします。
ステップ 4: ジョブコードの記述とジョブのデプロイ
-
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 つだけを残します。宛先テーブルに書き込む際、
METADATAtable_name を使用して、指定されたテーブルにデータをルーティングします。この方法では、AnalyticDB for PostgreSQL に作成するレプリケーションスロットは 1 つだけで済みます。これにより、ソースデータベースのリソース使用量が削減され、同期パフォーマンスが向上し、将来のメンテナンスが簡素化されます。 -
table-nameパラメーターを使用して複数のソーステーブルを指定します。テーブル名を括弧で囲み、縦棒 (|) で区切ります (例:(table1|table2|table3))。 -
debezium.snapshot.modeをneverに設定すると、ソーステーブルの増分データのみが同期されます。 完全なデータと増分データの両方を同期するには、設定をinitialに変更します。
-
-
ジョブ開発ページの上部で [Deep Check] をクリックして構文チェックを実行します。
-
[Deploy] をクリックし、次に [OK] をクリックします。
-
右上隅にある [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 を使用してフルデータをリアルタイムで読み書きする」をご参照ください。