このチュートリアルでは、リアルタイムエンリッチメントパイプラインの共通パターンである、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 以降が必要です。
代わりにカスタムコネクタを使用するには、「カスタムコネクタの管理」をご参照ください。
前提条件
開始する前に、次のものがあることを確認してください。
-
Flink 完全管理ワークスペース。「Flink 完全管理のアクティブ化」をご参照ください。
-
AnalyticDB for PostgreSQL インスタンスと特権アカウント。「インスタンスの作成」および「特権アカウントの作成」をご参照ください。
-
AnalyticDB for PostgreSQL インスタンスと Flink 完全管理ワークスペースが同じ VPC にあること。
異なる VPC にある場合は、「Flink 完全管理はどのようにして VPC をまたいでサービスにアクセスしますか?
ステップ 1:ホワイトリストの設定とデータの準備
-
AnalyticDB for PostgreSQL コンソールにログインします。
-
Flink 完全管理ワークスペースの CIDR ブロックを AnalyticDB for PostgreSQL インスタンスのホワイトリストに追加します。
-
ご利用の Flink 完全管理ワークスペースが使用する vSwitch の CIDR ブロックを検索します。「ホワイトリストの設定方法」をご参照ください。
-
CIDR ブロックを AnalyticDB for PostgreSQL インスタンスのホワイトリストに追加します。「手順」をご参照ください。
インターネット経由でインスタンスにアクセスする場合は、代わりにパブリック IP アドレスを追加してください。
-
-
インスタンス詳細ページで、右上の [データベースにログオン] をクリックし、ご利用のユーザー名とパスワードを入力します。詳細については、「クライアントツールを使用してインスタンスに接続する」をご参照ください。
-
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;を実行して、挿入されたデータを確認します。 -
Flink が出力を書き込むための
adbpg_sink_tableという名前の結果テーブルを作成します。CREATE TABLE adbpg_sink_table( id int, username text, score int );
ステップ 2:ストリームジョブの下書きの作成
-
Realtime Compute for Apache Flink コンソールにログインし、ご利用のワークスペースを見つけて、[操作] 列の [コンソール] をクリックします。
-
左側のナビゲーションウィンドウで、[開発] > [ETL] に移動します。SQL エディタページの左上隅で、[+] をクリックし、[新しい空白のストリームジョブの下書き] を選択します。
-
[新しい下書き] ダイアログボックスで、次のパラメーターを設定します。
パラメーター 説明 例 名前 下書きの名前。プロジェクト内で一意である必要があります。 adbpg-test場所 下書きが保存されるフォルダ。既存のフォルダの横にあるアイコンをクリックして、サブフォルダを作成します。 Draftエンジンバージョン Flink エンジンのバージョン。バージョンの詳細とライフサイクルについては、「エンジンバージョン」をご参照ください。 vvr-8.0.1-flink-1.17 -
[作成] をクリックします。
ステップ 3:下書きの作成とデプロイ
-
次の 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 に指示します。これは、各レコードが処理時に存在するディメンションデータでエンリッチ化され、後でディメンションテーブルが変更されても、すでに書き込まれた結果は更新されないことを意味します。 -
ディメンションテーブルと結果テーブルの接続パラメーターを更新します。
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いいえ — 書き込みリトライ間の待機時間 (ミリ秒)。 -
SQL エディタページの右上隅にある [検証] をクリックして、構文を確認します。
-
[デプロイ] をクリックします。
-
[運用保守] > [デプロイメント] ページで、ご利用のデプロイメントを見つけ、[操作] 列の [開始] をクリックします。
ステップ 4:結果の検証
-
AnalyticDB for PostgreSQL コンソールにログインします。
-
[データベースにログオン] をクリックします。詳細については、「クライアントからインスタンスに接続する」をご参照ください。
-
次のクエリを実行して、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 です。