Postgres CDC コネクタは、PostgreSQL データベースの完全なスナップショットを読み取り、exactly-once の処理セマンティクスで変更データをキャプチャします。
概要
Postgres CDC コネクタは、次の機能をサポートしています。
|
カテゴリ |
詳細 |
|
サポートされているタイプ |
SQL ソース、Flink CDC ソース 説明
シンクテーブルおよびルックアップ (ディメンション) テーブルには、JDBC コネクタを使用してください。 |
|
実行モード |
ストリーミング |
|
データフォーマット |
該当なし |
|
メトリクス |
|
|
API タイプ |
SQL および Flink CDC |
|
シンクの更新/削除 |
該当なし |
機能
VVR 8.0.6 以降、Postgres CDC コネクタはインクリメンタルスナップショットフレームワークと統合されています。履歴データ全体を読み取った後、自動的に WAL からの変更ログの読み取りに切り替わり、exactly-once セマンティクスを実現します。
主な機能は次のとおりです。
-
ストリーム処理およびバッチ処理の統合。単一のジョブで全データと増分データの両方を読み取ります。
-
スナップショットの並列読み取り。水平方向にスケールすることで、パフォーマンスを向上させます。
-
全データから増分データへのシームレスな切り替え。自動的にスケールインし、リソース使用量を削減します。
-
再開可能な読み取り。スナップショットフェーズ中にブレークポイントから再開することで、安定性を向上させます。
-
ロックフリーの読み取り。ロックが不要なため、オンラインオペレーションへの影響を回避できます。
前提条件
Postgres CDC コネクタは、PostgreSQL の論理レプリケーションを通じて CDC ストリームを読み取ります。ApsaraDB RDS for PostgreSQL、Amazon RDS for PostgreSQL、およびセルフマネージド PostgreSQL をサポートしています。
設定はデプロイメントのタイプによって異なります。Postgres の設定をご参照ください。
設定後、次の項目を確認してください。
-
wal_level が
logicalに設定されており、論理デコーディングが有効になっていること。 -
サブスクライブされた各テーブルの REPLICA IDENTITY が
FULLに設定されており、UPDATEおよびDELETEイベントにデータ整合性のための以前の列の値が含まれていること。説明REPLICA IDENTITYは、UPDATEおよびDELETEイベントに以前の列の値を含めるかどうかを制御する PostgreSQL のテーブルレベル設定です。詳細については、「REPLICA IDENTITY」をご参照ください。 -
max_wal_sendersおよびmax_replication_slotsの値が、使用中のスロット数と Flink ジョブに必要なスロット数の合計を超えていること。 -
アカウントに
SUPERUSER権限があるか、またはLOGINおよびREPLICATION権限の両方に加え、サブスクライブされたテーブルに対するSELECT権限があること。
-
Postgres テーブルに生成列が含まれている場合は、スロット作成時に publish_generated_columns パラメーターを
storedに設定してください。設定しない場合、スナップショットフェーズと増分フェーズのスキーマが異なる可能性があります。
注意事項
増分スナップショット機能には、VVR 8.0.6 以降が必要です。
レプリケーションスロット
Flink PostgreSQL CDC ジョブは、レプリケーションスロットを使用して、WAL の早期パージを防ぎ、データ整合性を確保します。スロットの管理が不適切な場合、過剰なディスク使用量や読み取り遅延を引き起こす可能性があります。ベストプラクティスは次のとおりです:
-
未使用スロットの速やかなクリーンアップ
-
Flink は、WAL データの損失を防ぐため、ジョブが停止したり、ステートレスで再起動したりした場合でも、レプリケーションスロットを自動的に削除しません。
-
ジョブを再開しない場合は、手動でそのレプリケーションスロットを削除してディスク領域を解放してください。
説明ライフサイクル管理:レプリケーションスロットをジョブレベルのリソースとして扱い、ジョブの開始と停止に合わせて管理してください。
-
-
古いスロットの再利用の回避
-
常に新しいスロット名を使用してください。古いスロットを再利用すると、ジョブは起動時に蓄積された履歴 WAL データを読み取らなければならなくなり、新しいデータの処理が遅延します。
-
PostgreSQL では、接続ごとに 1 つのスロットが必要です。各ジョブは一意のスロット名を使用する必要があります。
説明命名規則:
slot.nameをカスタマイズする場合、一時スロットとの競合を避けるため、my_slot_1のような数字のサフィックスを持つ名前は避けてください。
-
-
増分スナップショット有効時のスロットの動作
-
前提条件:チェックポイントが有効で、ソーステーブルにプライマリキーが定義されている必要があります。
-
スロット作成ルール:
-
増分スナップショットが無効な場合:並列度 1 のみがサポートされます。1 つのグローバルスロットが使用されます。
-
増分スナップショットが有効な場合:
-
スナップショットフェーズ:各並列ソースサブタスクは一時スロットを作成します。命名形式は
${slot.name}_${task_id}です。 -
増分フェーズ:すべての一時スロットは自動的に回収されます。1 つのグローバルスロットのみが保持されます。
-
-
-
最大スロット数:ソースの並列度 + 1 (スナップショットフェーズ中)
-
-
リソースとパフォーマンス
-
利用可能なスロットまたはディスク領域が限られている場合は、スナップショットの並列度を下げて、使用する一時スロットの数を減らしてください。これにより、スナップショットの読み取り速度が低下します。
-
ダウンストリームのシンクがべき等な書き込みをサポートしている場合は、
scan.incremental.snapshot.backfill.skip = trueを設定して、スナップショットフェーズ中の WAL のバックフィルをスキップし、起動を高速化します。これは at-least-once セマンティクスのみを提供し、必要な履歴の変更が失われる可能性があるため、ステートフルな計算 (集計やルックアップ結合) には適していません。
-
-
増分スナップショットが無効な場合、スナップショットフェーズ中にチェックポイントはサポートされません。
PostgreSQL パブリケーションの再利用
PostgreSQL CDC コネクタは、どのテーブルの変更をスロットにプッシュするかを決定するためにパブリケーションに依存しています。複数のジョブが同じパブリケーションを共有している場合、それらの設定は上書きされます。
原因
デフォルトの publication.autocreate.mode は filtered であり、コネクタ設定に含まれるテーブルのみが対象になります。これにより、ジョブの起動時にパブリケーションが変更され、他のジョブに影響を与える可能性があります。
解決策
-
監視対象のすべてのテーブルを含むパブリケーションを PostgreSQL で作成するか、ジョブごとに個別のパブリケーションを作成します。
-- my_flink_pub という名前でパブリケーションを作成し、すべてのテーブルを含める (または、指定したテーブルを対象に、ジョブごとにパブリケーションを作成) CREATE PUBLICATION my_flink_pub FOR TABLE table_a, table_b; -- または、より単純に、データベース内のすべてのテーブルを含める CREATE PUBLICATION my_flink_pub FOR ALL TABLES;説明すべてのテーブルをパブリケーションに含めることは、Flink クラスターの過剰な帯域幅と CPU 使用量が増加するため、大規模なデータベースでは推奨されません。
-
次の Flink 設定を追加します:
-
debezium.publication.name = 'my_flink_pub'(パブリケーション名を指定します) -
debezium.publication.autocreate.mode = 'disabled'(Flink が起動時にパブリケーションを作成または変更しようとするのを防ぎます)
-
これにより、完全なアイソレーションが実現され、新しいジョブが既存のジョブに影響を与えるのを防ぎます。
SQL
構文
CREATE TABLE postgrescdc_source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>',
'decoding.plugin.name'= 'pgoutput',
'scan.incremental.snapshot.enabled' = 'true',
-- バックフィルをスキップすると、読み取りが高速化され、リソース使用量を削減できますが、データの重複が発生する可能性があります。ダウンストリームシンクがべき等である場合に有効にしてください。
'scan.incremental.snapshot.backfill.skip' = 'false',
-- 本番環境では、この値を 'filtered' または 'disabled' に設定し、Flink 経由ではなく手動でパブリケーションを管理してください。
'debezium-publication.autocreate.mode' = 'disabled'
-- 複数のソースがある場合は、ソースごとに異なるパブリケーションを設定してください。
--'debezium.publication.name' = 'my_flink_pub'
);
コネクタオプション
|
オプション |
説明 |
データ型 |
必須 |
デフォルト |
備考 |
|
connector |
コネクタ名。 |
STRING |
はい |
– |
値は |
|
hostname |
PostgreSQL データベースの IP アドレスまたはホスト名。 |
STRING |
はい |
– |
– |
|
username |
PostgreSQL データベースサービスのユーザー名。 |
STRING |
はい |
– |
– |
|
password |
PostgreSQL データベースサービスのパスワード。 |
STRING |
はい |
– |
– |
|
database-name |
PostgreSQL データベース名。 |
STRING |
はい |
– |
データベースの名前。 |
|
schema-name |
PostgreSQL のスキーマ名。正規表現に対応しています。 |
STRING |
はい |
– |
スキーマ名は正規表現に対応しているため、複数のスキーマからデータを読み取ることができます。 |
|
table-name |
PostgreSQL のテーブル名。正規表現に対応しています。 |
STRING |
はい |
– |
テーブル名は正規表現に対応しているため、複数のテーブルからデータを読み取ることができます。 |
|
port |
ポート番号。 |
INTEGER |
いいえ |
5432 |
– |
|
decoding.plugin.name |
PostgreSQL の論理デコーディングプラグインの名前。 |
STRING |
いいえ |
decoderbufs |
この値は、PostgreSQL サービスにインストールされているプラグインによって決まります。対応しているプラグインは次のとおりです。
|
|
slot.name |
論理デコーディングスロット名。 |
STRING |
VVR 8.0.1 以降では必須。それ以前のバージョンではオプション。 |
|
VVR 8.0.1 以降ではデフォルト値はありません。 |
|
debezium.* |
Debezium のプロパティとパラメーター。 |
STRING |
いいえ |
– |
Debezium クライアントの動作をよりきめ細かく制御できます。 例: |
|
scan.incremental.snapshot.enabled |
インクリメンタルスナップショットを有効にするかどうかを指定します。 |
BOOLEAN |
いいえ |
false |
|
|
scan.startup.mode |
データ消費の起動モード。 |
STRING |
いいえ |
initial |
有効な値:
|
|
changelog-mode |
ストリームの変更をエンコードするためのチェンジログモード。 |
STRING |
いいえ |
all |
対応しているチェンジログモード:
|
|
heartbeat.interval.ms |
ハートビートパケットを送信する間隔。 |
Duration |
いいえ |
30000 |
単位はミリ秒です。 Postgres CDC コネクタは、スロットオフセットを進めるために、データベースに積極的にハートビートを送信します。テーブルの変更頻度が低い場合、この値を設定することで WAL ログの適時な回収が保証されます。 |
|
scan.incremental.snapshot.chunk.key-column |
スナップショットフェーズ中にシャードを分割するためのチャンクキーとして使用する列を指定します。 |
STRING |
いいえ |
– |
デフォルトでは、プライマリキーの最初の列が選択されます。 |
|
scan.incremental.close-idle-reader.enabled |
スナップショット完了後にアイドルリーダーを閉じるかどうかを指定します。 |
BOOLEAN |
いいえ |
false |
この設定を有効にするには、 |
|
scan.incremental.snapshot.backfill.skip |
スナップショットフェーズ中にログの読み取りをスキップするかどうかを指定します。 |
BOOLEAN |
いいえ |
false |
有効な値:
|
型マッピング
PostgreSQL から Flink への型マッピング
|
PostgreSQL CDC |
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 |
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 |
例
CREATE TABLE source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>'
);
SELECT * FROM source;
Flink CDC
VVR V11.4 以降では、PostgreSQL コネクタを Flink CDC ソースとしてサポートしています。
構文
source:
type: postgres
name: PostgreSQL Source
hostname: localhost
port: 5432
username: pg_username
password: pg_password
tables: db.scm.tbl
slot.name: test_slot
scan.startup.mode: initial
server-time-zone: UTC
connect.timeout: 120s
decoding.plugin.name: decoderbufs
sink:
type: ...
コネクタオプション
|
オプション |
説明 |
必須 |
データ型 |
デフォルト |
備考 |
|
type |
コネクタ名。 |
はい |
STRING |
– |
|
|
name |
データソース名。 |
いいえ |
STRING |
– |
– |
|
hostname |
PostgreSQL データベースサーバーのドメイン名または IP アドレス。 |
はい |
STRING |
– |
– |
|
port |
PostgreSQL データベースのポート。 |
いいえ |
INTEGER |
5432 |
– |
|
username |
PostgreSQL のユーザー名。 |
はい |
STRING |
– |
– |
|
password |
PostgreSQL のパスワード。 |
はい |
STRING |
– |
– |
|
tables |
キャプチャするテーブル名。 正規表現が使用できます。 |
はい |
STRING |
– |
重要
現在、同じデータベース内のテーブルのみをキャプチャできます。 ピリオド (.) は完全修飾名の区切り文字として扱われます。正規表現で任意の文字に一致させるためにピリオド (.) を使用する場合は、バックスラッシュでエスケープしてください。例: |
|
slot.name |
PostgreSQL レプリケーションスロット名。 |
はい |
STRING |
– |
名前は PostgreSQL レプリケーションスロットの命名規則に準拠し、小文字、数字、アンダースコアを含めることができます。 |
|
decoding.plugin.name |
サーバーにインストールされている PostgreSQL 論理デコーディングプラグインの名前。 |
いいえ |
STRING |
|
有効な値: |
|
tables.exclude |
除外するテーブル。このオプションは |
いいえ |
STRING |
– |
|
|
server-time-zone |
データベースサーバーのセッションタイムゾーン (例: "Asia/Shanghai")。 |
いいえ |
STRING |
– |
設定されていない場合、システムのデフォルトタイムゾーン ( |
|
scan.incremental.snapshot.chunk.size |
インクリメンタルスナップショットフレームワークにおける各チャンクのサイズ (行数)。 |
いいえ |
INTEGER |
8096 |
インクリメンタルスナップショットが有効な場合、テーブルは読み取り用に複数のチャンクに分割されます。チャンクのデータは、完全に消費される前にメモリにキャッシュされます。 チャンクを小さくすると、テーブルの総チャンク数が増加します。これにより障害復旧の粒度は小さくなりますが、メモリ不足 (OOM) エラーが発生し、全体的なスループットが低下する可能性があります。そのため、バランスを取り、適切なチャンクサイズを設定する必要があります。 |
|
scan.snapshot.fetch.size |
テーブルの全データを読み取る際に一度に取得するレコードの最大数。 |
いいえ |
INTEGER |
1024 |
– |
|
scan.startup.mode |
データ消費の起動モード。 |
いいえ |
STRING |
initial |
有効な値:
|
|
scan.incremental.close-idle-reader.enabled |
スナップショット完了後にアイドルリーダーを閉じるかどうかを指定します。 |
いいえ |
BOOLEAN |
false |
この設定を有効にするには、 |
|
scan.lsn-commit.checkpoints-num-delay |
LSN オフセットのコミットを開始する前に遅延させるチェックポイントの数。 |
いいえ |
INTEGER |
3 |
チェックポイント LSN オフセットはローリング方式でコミットされ、状態から復旧できなくなることを防ぎます。 |
|
connect.timeout |
コネクタが PostgreSQL データベースサーバーへの接続を試行する際の、タイムアウトまでの最大待機時間。 |
いいえ |
DURATION |
30s |
この値は 250 ミリ秒未満にすることはできません。 |
|
connect.max-retries |
コネクタが接続を確立するための最大再試行回数。 |
いいえ |
INTEGER |
3 |
– |
|
connection.pool.size |
コネクションプールのサイズ。 |
いいえ |
INTEGER |
20 |
– |
|
jdbc.properties.* |
ユーザーがカスタム JDBC URL プロパティを渡すことができます。 |
いいえ |
STRING |
– |
ユーザーは、 |
|
heartbeat.interval |
最新の利用可能な WAL ログオフセットを追跡するためのハートビートイベントの送信間隔。 |
いいえ |
DURATION |
30s |
– |
|
debezium.* |
PostgreSQL サーバーからのデータ変更をキャプチャするために使用される Debezium Embedded Engine に Debezium プロパティを渡します。 |
いいえ |
STRING |
– |
Debezium PostgreSQL コネクタのプロパティについては、Debezium ドキュメントをご参照ください。 |
|
chunk-meta.group.size |
チャンクメタデータのサイズ。 |
いいえ |
INTEGER |
1000 |
メタデータがこの値より大きい場合、分割して渡されます。 |
|
metadata.list |
ダウンストリームに渡される利用可能なメタデータのリスト。transform モジュールで使用できます。 |
いいえ |
STRING |
– |
区切り文字としてカンマ (,) を使用します。現在利用可能なメタデータは |
|
scan.incremental.snapshot.unbounded-chunk-first.enabled |
スナップショット読み取りフェーズ中に、無制限チャンクを最初にディスパッチするかどうかを指定します。 |
いいえ |
BOOLEAN |
false |
これは実験的機能です。有効にすると、スナップショットフェーズ中に TaskManager が最後のチャンクを同期する際の OOM エラーのリスクを軽減できます。ジョブの初回起動前にこれを追加することを推奨します。 |