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

Realtime Compute for Apache Flink:AnalyticDB for PostgreSQL へのアクセス

最終更新日:Jul 25, 2026

このチュートリアルでは、リアルタイムエンリッチメントパイプラインの共通パターンである、Flink SQL ジョブで AnalyticDB for PostgreSQL をディメンションテーブルと結果テーブルの両方として使用する方法を説明します。

このチュートリアルを完了すると、Datagen ソースからデータを読み取り、AnalyticDB for PostgreSQL ディメンションテーブルからユーザー情報をルックアップし、エンリッチ化されたレコードを AnalyticDB for PostgreSQL 結果テーブルに書き込む Flink ジョブが実行されるようになります。

制限事項

  • Realtime Compute for Apache Flink は、サーバーレスモードの AnalyticDB for PostgreSQL からデータを読み取ることはできません。

  • AnalyticDB for PostgreSQL コネクタには、Ververica Runtime (VVR) 6.0.0 以降が必要です。

  • AnalyticDB for PostgreSQL V7.0 には、VVR 8.0.1 以降が必要です。

代わりにカスタムコネクタを使用するには、「カスタムコネクタの管理」をご参照ください。

前提条件

開始する前に、次のものがあることを確認してください。

異なる VPC にある場合は、「Flink 完全管理はどのようにして VPC をまたいでサービスにアクセスしますか?

ステップ 1:ホワイトリストの設定とデータの準備

  1. AnalyticDB for PostgreSQL コンソールにログインします。

  2. Flink 完全管理ワークスペースの CIDR ブロックを AnalyticDB for PostgreSQL インスタンスのホワイトリストに追加します。

    1. ご利用の Flink 完全管理ワークスペースが使用する vSwitch の CIDR ブロックを検索します。「ホワイトリストの設定方法」をご参照ください。

    2. CIDR ブロックを AnalyticDB for PostgreSQL インスタンスのホワイトリストに追加します。「手順」をご参照ください。

    インターネット経由でインスタンスにアクセスする場合は、代わりにパブリック IP アドレスを追加してください。
  3. インスタンス詳細ページで、右上の [データベースにログオン] をクリックし、ご利用のユーザー名とパスワードを入力します。詳細については、「クライアントツールを使用してインスタンスに接続する」をご参照ください。

  4. adbpg_dim_table という名前のディメンションテーブルを作成し、50 行のサンプルデータを挿入します。

    -- ディメンションテーブルの作成
    CREATE TABLE adbpg_dim_table(
      id int,
      username text,
      PRIMARY KEY(id)
    );
    
    -- 50 行を挿入:id は 1 から 50 の範囲、username は "username" に行番号を続けたもの
    INSERT INTO adbpg_dim_table(id, username)
    SELECT i, 'username'||i::text
    FROM generate_series(1, 50) AS t(i);

    SELECT * FROM adbpg_dim_table ORDER BY id; を実行して、挿入されたデータを確認します。

  5. Flink が出力を書き込むための adbpg_sink_table という名前の結果テーブルを作成します。

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

ステップ 2:ストリームジョブの下書きの作成

  1. Realtime Compute for Apache Flink コンソールにログインし、ご利用のワークスペースを見つけて、[操作] 列の [コンソール] をクリックします。

  2. 左側のナビゲーションウィンドウで、[開発] > [ETL] に移動します。SQL エディタページの左上隅で、[+] をクリックし、[新しい空白のストリームジョブの下書き] を選択します。

  3. [新しい下書き] ダイアログボックスで、次のパラメーターを設定します。

    パラメーター 説明
    名前 下書きの名前。プロジェクト内で一意である必要があります。 adbpg-test
    場所 下書きが保存されるフォルダ。既存のフォルダの横にあるアイコンをクリックして、サブフォルダを作成します。 Draft
    エンジンバージョン Flink エンジンのバージョン。バージョンの詳細とライフサイクルについては、「エンジンバージョン」をご参照ください。 vvr-8.0.1-flink-1.17
  4. [作成] をクリックします。

ステップ 3:下書きの作成とデプロイ

  1. 次の SQL をコードエディタにコピーします。この SQL は、3 つのテーブルと、Datagen ストリームを AnalyticDB for PostgreSQL のユーザーデータでエンリッチ化するルックアップ結合を定義します。

    -- ソーステーブル:Datagen はシーケンシャルな ID (1-50) とランダムなスコア (70-100) を生成します。
    -- この例では、WITH 句の変更は不要です。
    CREATE TEMPORARY TABLE datagen_source (
      id INT,
      score INT
    ) WITH (
      'connector' = 'datagen',
      'fields.id.kind' = 'sequence',
      'fields.id.start' = '1',
      'fields.id.end' = '50',
      'fields.score.kind' = 'random',
      'fields.score.min' = '70',
      'fields.score.max' = '100'
    );
    
    -- ディメンションテーブル:AnalyticDB for PostgreSQL によってバックアップされます。
    -- Flink は処理時にこのテーブルをクエリして、ID によってユーザー名をルックアップします。
    -- WITH 句の値を実際の接続詳細に置き換えてください。
    CREATE TEMPORARY TABLE dim_adbpg(
      id int,
      username varchar,
      PRIMARY KEY(id) NOT ENFORCED
    ) WITH (
      'connector' = 'adbpg',
      'url' = 'jdbc:postgresql://gp-2ze****3tysk255b5-master.gpdb.rds.aliyuncs.com:5432/flinktest',
      'tablename' = 'adbpg_dim_table',
      'username' = 'flinktest',
      'password' = '${secret_values.adb_password}',
      'maxRetryTimes' = '2',
      'cache' = 'lru',
      'cacheSize' = '100'
    );
    
    -- 結果テーブル:Flink はエンリッチ化されたレコードをここに書き込みます。
    -- WITH 句の値を実際の接続詳細に置き換えてください。
    CREATE TEMPORARY TABLE sink_adbpg (
      id int,
      username varchar,
      score int
    ) WITH (
      'connector' = 'adbpg',
      'url' = 'jdbc:postgresql://gp-2ze****3tysk255b5-master.gpdb.rds.aliyuncs.com:5432/flinktest',
      'tablename' = 'adbpg_sink_table',
      'username' = 'flinktest',
      'password' = '${secret_values.adb_password}',
      'maxRetryTimes' = '2',
      'conflictMode' = 'ignore',
      'retryWaitTime' = '200'
    );
    
    -- ルックアップ結合:datagen_source からの各レコードについて、Flink はレコードが処理された時点 (PROCTIME()) で dim_adbpg 内の一致する行をルックアップします。
    INSERT INTO sink_adbpg
    SELECT ts.id, ts.username, ds.score
    FROM datagen_source AS ds
    JOIN dim_adbpg FOR SYSTEM_TIME AS OF PROCTIME() AS ts
    ON ds.id = ts.id;

    ルックアップ結合の構文について: FOR SYSTEM_TIME AS OF PROCTIME() は、各ソースレコードが処理された時点でディメンションテーブルをルックアップするように Flink に指示します。これは、各レコードが処理時に存在するディメンションデータでエンリッチ化され、後でディメンションテーブルが変更されても、すでに書き込まれた結果は更新されないことを意味します。

  2. ディメンションテーブルと結果テーブルの接続パラメーターを更新します。WITH 句のプレースホルダー値を、ご利用の実際の AnalyticDB for PostgreSQL 接続詳細に置き換えます。Datagen ソーステーブルは変更する必要はありません。完全なパラメーターリファレンスとデータ型マッピングについては、「AnalyticDB for PostgreSQL コネクタ」をご参照ください。

    パラメーター 必須 デフォルト 説明
    url はい フォーマット jdbc:postgresql://<内部エンドポイント>:<ポート>/<データベース名> の JDBC URL。これは、AnalyticDB for PostgreSQL コンソールのインスタンスの [データベース接続] ページで確認できます。
    tablename はい AnalyticDB for PostgreSQL データベース内のテーブル名。
    username はい データベースにアクセスするためのユーザー名。
    password はい データベースアカウントのパスワード。
    targetSchema いいえ public スキーマ名。テーブルが public スキーマにない場合にのみ指定します。
    maxRetryTimes いいえ 書き込み失敗後の最大リトライ回数。
    cache いいえ ディメンションテーブルのルックアップのキャッシュポリシー。lru に設定すると、最近アクセスされたエントリがメモリに保持されます。LRU キャッシュはデータベースのトラフィックを削減し、ルックアップのスループットを向上させますが、キャッシュされたエントリは古くなる可能性があります。これはスループットとデータの鮮度のトレードオフです。有効にする前に、cacheSize を調整し、古いデータに対する許容度を考慮してください。
    cacheSize いいえ キャッシュするエントリの最大数。値を大きくすると、データベースリクエストは減少しますが、より多くのメモリを消費します。
    conflictMode いいえ 書き込みが既存のプライマリキーまたはインデックスと競合した場合に実行されるアクション。ignore に設定すると、競合する行をスキップします。
    retryWaitTime いいえ 書き込みリトライ間の待機時間 (ミリ秒)。
  3. SQL エディタページの右上隅にある [検証] をクリックして、構文を確認します。

  4. [デプロイ] をクリックします。

  5. [運用保守] > [デプロイメント] ページで、ご利用のデプロイメントを見つけ、[操作] 列の [開始] をクリックします。

ステップ 4:結果の検証

  1. AnalyticDB for PostgreSQL コンソールにログインします。

  2. [データベースにログオン] をクリックします。詳細については、「クライアントからインスタンスに接続する」をご参照ください。

  3. 次のクエリを実行して、Flink が結果テーブルに書き込んだレコードを表示します。

    SELECT * FROM adbpg_sink_table ORDER BY id;

    結果には 50 行が含まれ、各行にはユーザー ID、ディメンションテーブルからの一致するユーザー名、および 70 から 100 の間のランダムなスコアが含まれている必要があります。

    このクエリは、id、username、score の 3 つの列を持つ 7 つのレコードを返します。行 1 から 7 は username1 から username7 に対応し、スコアはそれぞれ 94、79、70、93、71、82、87 です。

参考