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

Realtime Compute for Apache Flink:PostgreSQL CDC

最終更新日:Aug 21, 2026

Postgres CDC コネクタは、PostgreSQL データベースの完全なスナップショットを読み取り、exactly-once の処理セマンティクスで変更データをキャプチャします。

概要

Postgres CDC コネクタは、次の機能をサポートしています。

カテゴリ

詳細

サポートされているタイプ

SQL ソース、Flink CDC ソース

説明

シンクテーブルおよびルックアップ (ディメンション) テーブルには、JDBC コネクタを使用してください。

実行モード

ストリーミング

データフォーマット

該当なし

メトリクス

監視メトリクス

  • currentFetchEventTimeLag :データが生成されてから、ソースオペレーターが取得するまでの間隔。

  • currentEmitEventTimeLag :データが生成されてから、ソースオペレーターから送出するまでの間隔。

  • sourceIdleTime :ソースが新しいデータを生成していない期間。

説明
  • currentFetchEventTimeLag および currentEmitEventTimeLag メトリクスは、増分フェーズでのみ有効です。スナップショットフェーズでは、これらの値は常に 0 になります。

  • メトリクスの詳細については、「メトリクスの説明」をご参照ください。

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 セマンティクスのみを提供し、必要な履歴の変更が失われる可能性があるため、ステートフルな計算 (集計やルックアップ結合) には適していません。

  • 増分スナップショットが無効な場合、スナップショットフェーズ中にチェックポイントはサポートされません。

    スナップショットフェーズ中のタイムアウトを回避するための設定

    増分スナップショットが無効な場合、スナップショットフェーズ中のチェックポイントが原因でタイムアウトによるフェイルオーバーが発生する可能性があります。これらのパラメーターを [Other Configuration] (カスタム実行パラメーターの設定) で設定してください:

    execution.checkpointing.interval: 10min
    execution.checkpointing.tolerable-failed-checkpoints: 100
    restart-strategy: fixed-delay
    restart-strategy.fixed-delay.attempts: 2147483647

    パラメーター:

    パラメーター

    説明

    備考

    execution.checkpointing.interval

    チェックポイント間の間隔。

    単位は 10min や 30s のような期間形式です。

    execution.checkpointing.tolerable-failed-checkpoints

    ジョブが失敗するまでに許容されるチェックポイントの失敗回数。

    このパラメーターとチェックポイントスケジューリング間隔の積が、許容されるスナップショット読み取り時間となります。

    説明

    テーブルが非常に大きい場合は、このパラメーターに大きな値を設定してください。

    restart-strategy

    ジョブの再起動戦略。

    有効な値:

    • fixed-delay:固定遅延再起動戦略。

    • failure-rate:失敗率再起動戦略。

    • exponential-delay:指数遅延再起動戦略。

    詳細については、「再起動戦略」をご参照ください。

    restart-strategy.fixed-delay.attempts

    fixed-delay 再起動戦略の最大再起動試行回数。

    –

PostgreSQL パブリケーションの再利用

PostgreSQL CDC コネクタは、どのテーブルの変更をスロットにプッシュするかを決定するためにパブリケーションに依存しています。複数のジョブが同じパブリケーションを共有している場合、それらの設定は上書きされます。

原因

デフォルトの publication.autocreate.mode は filtered であり、コネクタ設定に含まれるテーブルのみが対象になります。これにより、ジョブの起動時にパブリケーションが変更され、他のジョブに影響を与える可能性があります。

解決策

  1. 監視対象のすべてのテーブルを含むパブリケーションを 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 使用量が増加するため、大規模なデータベースでは推奨されません。

  2. 次の 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

はい

–

値は postgres-cdc にする必要があります。

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 サービスにインストールされているプラグインによって決まります。対応しているプラグインは次のとおりです。

  • decoderbufs:PostgreSQL 9.6 以降に対応しています。このプラグインのインストールが必要です。

  • pgoutput (推奨):PostgreSQL 10 以降の公式組み込みプラグイン。

slot.name

論理デコーディングスロット名。

STRING

VVR 8.0.1 以降では必須。それ以前のバージョンではオプション。

flink (8.0.1 より前)

PSQLException: ERROR: replication slot "debezium" is active for PID 974 エラーを回避するために、各テーブルに一意の slot.name を設定してください。レプリケーションスロット。

VVR 8.0.1 以降ではデフォルト値はありません。

debezium.*

Debezium のプロパティとパラメーター。

STRING

いいえ

–

Debezium クライアントの動作をよりきめ細かく制御できます。 例: 'debezium.snapshot.mode' = 'never'。 設定プロパティ。

scan.incremental.snapshot.enabled

インクリメンタルスナップショットを有効にするかどうかを指定します。

BOOLEAN

いいえ

false

説明
  • これは VVR 8.0.6 以降で利用可能な実験的機能です。

  • 利点、前提条件、および制限事項については、「機能」、「前提条件」、および「使用上の注意」をご参照ください。

scan.startup.mode

データ消費の起動モード。

STRING

いいえ

initial

有効な値:

  • initial:初回起動時にすべての履歴データをスキャンし、その後最新の WAL データを読み取ります。

  • latest-offset:初回起動時にすべての履歴データをスキャンしません。WAL の末尾から読み取りを開始するため、コネクタの起動後に行われた最新の変更のみを読み取ります。

  • snapshot:すべての履歴データをスキャンし、スナップショットフェーズ中に生成された新しい WAL データを読み取った後、ジョブが停止します。

changelog-mode

ストリームの変更をエンコードするためのチェンジログモード。

STRING

いいえ

all

対応しているチェンジログモード:

  • ALL:INSERT、DELETE、UPDATE_BEFORE、UPDATE_AFTER を含むすべての種類に対応します。

  • UPSERT:INSERT、DELETE、UPDATE_AFTER を含む upsert タイプのみに対応します。

heartbeat.interval.ms

ハートビートパケットを送信する間隔。

Duration

いいえ

30000

単位はミリ秒です。

Postgres CDC コネクタは、スロットオフセットを進めるために、データベースに積極的にハートビートを送信します。テーブルの変更頻度が低い場合、この値を設定することで WAL ログの適時な回収が保証されます。

scan.incremental.snapshot.chunk.key-column

スナップショットフェーズ中にシャードを分割するためのチャンクキーとして使用する列を指定します。

STRING

いいえ

–

デフォルトでは、プライマリキーの最初の列が選択されます。

scan.incremental.close-idle-reader.enabled

スナップショット完了後にアイドルリーダーを閉じるかどうかを指定します。

BOOLEAN

いいえ

false

この設定を有効にするには、execution.checkpointing.checkpoints-after-tasks-finish.enabled を true に設定してください。

scan.incremental.snapshot.backfill.skip

スナップショットフェーズ中にログの読み取りをスキップするかどうかを指定します。

BOOLEAN

いいえ

false

有効な値:

  • true:スキップします。

    増分フェーズでは、下限ウォーターマークからログの読み取りが開始されます。

    ダウンストリームのオペレーターまたはストレージがべき等性に対応している場合は、フルフェーズでのログ読み取りをスキップすることを推奨します。これにより WAL スロット数が削減されますが、at-least-once セマンティクスのみが保証されます。

  • false:スキップしません。

    フルフェーズで分割を読み取る際、下限ウォーターマークと上限ウォーターマークの間のログが読み取られ、整合性が保証されます。

    SQL で集計、結合、または類似の操作を実行する場合は、フルフェーズでのログ読み取りをスキップしないことを推奨します。

型マッピング

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

–

postgres である必要があります。

name

データソース名。

いいえ

STRING

–

–

hostname

PostgreSQL データベースサーバーのドメイン名または IP アドレス。

はい

STRING

–

–

port

PostgreSQL データベースのポート。

いいえ

INTEGER

5432

–

username

PostgreSQL のユーザー名。

はい

STRING

–

–

password

PostgreSQL のパスワード。

はい

STRING

–

–

tables

キャプチャするテーブル名。

正規表現が使用できます。

はい

STRING

–

重要

現在、同じデータベース内のテーブルのみをキャプチャできます。

ピリオド (.) は完全修飾名の区切り文字として扱われます。正規表現で任意の文字に一致させるためにピリオド (.) を使用する場合は、バックスラッシュでエスケープしてください。例: bdb.schema_\.*.order_\.*

slot.name

PostgreSQL レプリケーションスロット名。

はい

STRING

–

名前は PostgreSQL レプリケーションスロットの命名規則に準拠し、小文字、数字、アンダースコアを含めることができます。

decoding.plugin.name

サーバーにインストールされている PostgreSQL 論理デコーディングプラグインの名前。

いいえ

STRING

pgoutput

有効な値: decoderbufs、pgoutput

tables.exclude

除外するテーブル。このオプションは tables オプションの後に適用されます。正規表現が使用できます。

いいえ

STRING

–

tables オプションをご参照ください。

server-time-zone

データベースサーバーのセッションタイムゾーン (例: "Asia/Shanghai")。

いいえ

STRING

–

設定されていない場合、システムのデフォルトタイムゾーン (ZoneId.systemDefault()) が使用されます。

scan.incremental.snapshot.chunk.size

インクリメンタルスナップショットフレームワークにおける各チャンクのサイズ (行数)。

いいえ

INTEGER

8096

インクリメンタルスナップショットが有効な場合、テーブルは読み取り用に複数のチャンクに分割されます。チャンクのデータは、完全に消費される前にメモリにキャッシュされます。

チャンクを小さくすると、テーブルの総チャンク数が増加します。これにより障害復旧の粒度は小さくなりますが、メモリ不足 (OOM) エラーが発生し、全体的なスループットが低下する可能性があります。そのため、バランスを取り、適切なチャンクサイズを設定する必要があります。

scan.snapshot.fetch.size

テーブルの全データを読み取る際に一度に取得するレコードの最大数。

いいえ

INTEGER

1024

–

scan.startup.mode

データ消費の起動モード。

いいえ

STRING

initial

有効な値:

  • initial (デフォルト): 初回起動時にスナップショットをスキャンし、その後最新の WAL データに切り替えます。

  • latest-offset: スナップショットの読み取りをスキップします。WAL の末尾から読み取りを開始するため、コネクタの起動後に行われた最新の変更のみを読み取ります。

  • committed-offset: スナップショットの読み取りをスキップします。指定されたオフセットから WAL データを消費します。

  • snapshot: スナップショットのみを消費し、増分データは消費しません。

scan.incremental.close-idle-reader.enabled

スナップショット完了後にアイドルリーダーを閉じるかどうかを指定します。

いいえ

BOOLEAN

false

この設定を有効にするには、execution.checkpointing.checkpoints-after-tasks-finish.enabled を true に設定してください。

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

–

ユーザーは、'jdbc.properties.useSSL' = 'false' などのカスタムプロパティを渡すことができます。

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

–

区切り文字としてカンマ (,) を使用します。現在利用可能なメタデータは op_ts です。

scan.incremental.snapshot.unbounded-chunk-first.enabled

スナップショット読み取りフェーズ中に、無制限チャンクを最初にディスパッチするかどうかを指定します。

いいえ

BOOLEAN

false

これは実験的機能です。有効にすると、スナップショットフェーズ中に TaskManager が最後のチャンクを同期する際の OOM エラーのリスクを軽減できます。ジョブの初回起動前にこれを追加することを推奨します。

リファレンス