このトピックでは、Realtime Compute for Apache Flink を使用して AnalyticDB for PostgreSQL にデータを書き込む方法について説明します。
制限事項
-
この機能は、AnalyticDB for PostgreSQL の Serverless モードをサポートしていません。
-
AnalyticDB for PostgreSQL コネクタは、Realtime Compute for Apache Flink VVR 6.0.0 以降でのみサポートされています。
-
AnalyticDB for PostgreSQL V7.0 は、Realtime Compute for Apache Flink VVR 8.0.1 以降でのみサポートされています。
説明カスタムコネクタを使用する場合、詳細については、「カスタムコネクタの管理」をご参照ください。
前提条件
-
フルマネージド Flink のワークスペースが作成済みであること。 詳細については、「Realtime Compute for Apache Flink の有効化」をご参照ください。
-
AnalyticDB for PostgreSQL インスタンスが作成済みであること。 詳細については、「インスタンスの作成」をご参照ください。
-
AnalyticDB for PostgreSQL インスタンスとフルマネージド Flink のワークスペースは、同じ VPC にデプロイする必要があります。
AnalyticDB for PostgreSQL インスタンスの設定
-
AnalyticDB for PostgreSQL コンソールにログインします。
-
Flink ワークスペースの CIDR ブロックを AnalyticDB for PostgreSQL インスタンスのホワイトリストに追加します。 ホワイトリストの設定方法の詳細については、「ホワイトリストの設定」をご参照ください。
-
[データベースにログオン] をクリックします。データベースに接続する他の方法の詳細については、「クライアントを使用してインスタンスに接続する」をご参照ください。
-
AnalyticDB for PostgreSQL インスタンスでテーブルを作成します。
テーブルを作成する SQL ステートメントの例:
CREATE TABLE test_adbpg_table( b1 int, b2 int, b3 text, PRIMARY KEY(b1) );
Realtime Compute for Apache Flink の設定
-
Realtime Compute コンソールにログインします。
-
[Fully Managed Flink] タブで、対象のワークスペースの [アクション] 列にある [コンソール] をクリックします。
-
左側のナビゲーションペインで、[コネクタ] をクリックします。
-
[コネクタ] ページで、[カスタムコネクタの作成] をクリックします。
-
カスタムコネクタの JAR ファイルをアップロードします。
説明-
カスタムの AnalyticDB for PostgreSQL Flink コネクタの JAR パッケージを取得します。 詳細については、「AnalyticDB PostgreSQL Connector」をご参照ください。
-
JAR パッケージのバージョンを、Realtime Compute for Apache Flink のエンジンバージョンと一致させてください。
-
-
アップロードが完了したら、[次へ] をクリックします。
システムがアップロードされたカスタムコネクタを解析します。 解析が成功した場合、次のステップに進むことができます。 解析に失敗した場合は、アップロードしたカスタムコネクタのコードが Apache Flink コミュニティの標準に準拠しているか確認してください。
-
[完了] をクリックします。
作成されたカスタムコネクタがコネクタ一覧に表示されます。
Flink ジョブの作成
-
Realtime Compute コンソールにログインします。Fully Managed Flink タブで、対象のワークスペースの [アクション] 列の [コンソール] をクリックします。
-
左側のナビゲーションペインで、[SQL 開発] をクリックします。[新規] をクリックし、[空白のストリーミングジョブドラフト] を選択して、[次へ] をクリックします。
-
[新規ファイルドラフト] ダイアログボックスで、ジョブパラメーターを設定します。
ジョブパラメーター
説明
例
[名前]
ジョブの名前。
説明ジョブ名は、現在のプロジェクト内で一意にする必要があります。
adbpg-test
[場所]
ジョブのコードファイルが属するフォルダー。
既存のフォルダーの右側にある
アイコンをクリックして、サブフォルダーを作成することもできます。Job Drafts
[エンジンバージョン]
現在のジョブで使用する Flink エンジンバージョン。 エンジンバージョン番号、バージョンマッピング、ライフサイクルマイルストーンの詳細については、「エンジンバージョン」をご参照ください。
vvr-6.0.7-flink-1.15
-
[作成] をクリックします。
AnalyticDB for PostgreSQL へのデータ書き込み
-
ジョブコードを記述します。
AnalyticDB for PostgreSQL に、ランダムなソーステーブル
datagen_sourceと宛先テーブルtest_adbpg_tableを作成します。次のジョブコードをジョブテキストエディターにコピーします。CREATE TABLE datagen_source ( f_sequence INT, f_random INT, f_random_str STRING ) WITH ( 'connector' = 'datagen', 'rows-per-second'='5', 'fields.f_sequence.kind'='sequence', 'fields.f_sequence.start'='1', 'fields.f_sequence.end'='1000', 'fields.f_random.min'='1', 'fields.f_random.max'='1000', 'fields.f_random_str.length'='10' ); CREATE TABLE test_adbpg_table ( `B1` bigint, `B2` bigint, `B3` VARCHAR, PRIMARY KEY(B1) not ENFORCED ) with ( 'connector' = 'adbpg-nightly-1.13', 'password' = 'xxx', 'tablename' = 'test_adbpg_table', 'username' = 'xxxx', 'url' = 'jdbc:postgresql://url:5432/schema', 'maxretrytimes' = '2', 'batchsize' = '50000', 'connectionmaxactive' = '5', 'conflictmode' = 'ignore', 'usecopy' = '0', 'targetschema' = 'public', 'exceptionmode' = 'ignore', 'casesensitive' = '0', 'writemode' = '1', 'retrywaittime' = '200' );datagen_sourceテーブルのパラメーターを変更する必要はありません。test_adbpg_tableテーブルのパラメーターは、実際の業務要件に基づいて変更する必要があります。以下の表でパラメーターについて説明します。パラメーター
必須
説明
connector
はい
コネクターの名前です。このパラメーターを
adbpg-nightly-<version number>に設定します。例:adbpg-nightly-1.13。url
はい
AnalyticDB for PostgreSQL インスタンスの JDBC URL。 形式:
jdbc:postgresql://<内部エンドポイント>:<ポート>/<データベース名>。 例:jdbc:postgresql://gp-xxxxxx.gpdb.cn-chengdu.rds.aliyuncs.com:5432/postgres。tablename
はい
AnalyticDB for PostgreSQL テーブルの名前。
username
はい
AnalyticDB for PostgreSQL インスタンスのデータベースアカウント。
password
はい
AnalyticDB for PostgreSQL インスタンスのデータベースアカウントのパスワード。
maxretrytimes
いいえ
SQL の実行に失敗した場合の最大リトライ回数。 デフォルト値:3。
batchsize
いいえ
1 回のバッチで書き込むデータレコードの最大数。 デフォルト値:50000。
exceptionmode
いいえ
データ書き込み中に例外が発生した場合のエラー処理ポリシー。 有効な値:
-
ignore:例外の原因となったデータを無視します。 これがデフォルト値です。
-
strict:データの書き込み中に例外が発生した場合、フェールオーバーをトリガーしてエラーを報告します。
conflictmode
いいえ
プライマリキーまたは一意なインデックスの競合を処理するためのポリシー。 有効な値:
-
ignore:プライマリキーの競合を無視し、既存のデータを保持します。
-
strict:プライマリキーの競合が発生した場合、フェールオーバーをトリガーしてエラーを報告します。
-
update:プライマリキーの競合が発生した場合にデータを更新します。
-
upsert:プライマリキーの競合が発生した場合、UPSERT メソッドを使用してデータを書き込みます。 これがデフォルト値です。
AnalyticDB for PostgreSQL は、INSERT ON CONFLICT と COPY ON CONFLICT を使用して UPSERT を実装します。 宛先テーブルがパーティションテーブルの場合、マイナーカーネルバージョンは V6.3.6.1 以降である必要があります。 マイナーカーネルバージョンのアップグレード方法の詳細については、「エンジンバージョンのアップグレード」をご参照ください。
targetschema
いいえ
AnalyticDB for PostgreSQL データベースのスキーマ。 デフォルト値:public。
writemode
いいえ
データ書き込みモード。 有効な値:
-
0:BATCH INSERT を使用してデータを書き込みます。
-
1:COPY API を使用してデータを書き込みます。 これがデフォルト値です。
-
2:BATCH UPSERT を使用してデータを書き込みます。
verbose
いいえ
コネクタのランタイムログを出力するかどうかを指定します。 有効な値:
-
0:ランタイムログを出力しません。 これがデフォルト値です。
-
1:ランタイムログを出力します。
retrywaittime
いいえ
例外発生時のリトライ間隔。 単位:ミリ秒。 デフォルト値:100。
batchwritetimeoutms
いいえ
バッチ書き込みにおける最大バッチ蓄積時間。 この時間を超えると、蓄積されたバッチが書き込まれます。 単位:ミリ秒。 デフォルト値:50000。
connectionmaxactive
いいえ
接続プールのパラメーター。 このパラメーターは、単一のタスクマネージャーの接続プールにおける同時接続数の最大値を指定します。 デフォルト値:5。
casesensitive
いいえ
列名とテーブル名で大文字と小文字を区別するかどうかを指定します。 有効な値:
-
0:大文字と小文字を区別しません。 これがデフォルト値です。
-
1:大文字と小文字を区別します。
説明パラメーターと型のマッピングをサポートしています。 詳細については、AnalyticDB for PostgreSQL (ADB PG) のコネクタドキュメントをご参照ください。
-
-
ジョブを開始します。
-
ジョブ開発ページの上部で、[デプロイ] をクリックします。表示されるダイアログボックスで、[OK] をクリックします。
説明セッションクラスターは、非本番環境での開発とテストに適しています。 セッションクラスターを使用して、ジョブのデバッグ、ジョブマネージャー (JM) のリソース使用率の向上、ジョブの起動の高速化が可能です。 ただし、ビジネスの安定性に関する懸念から、ジョブをセッションクラスターに送信することは推奨しません。 詳細については、「ジョブのデバッグ」をご参照ください。
-
[デプロイメント] ページで、対象のジョブの [アクション] 列で [再開] をクリックします。
-
[再開] をクリックします。
-
結果の検証
-
AnalyticDB for PostgreSQL データベースに接続します。 詳細については、「クライアントを使用したインスタンスへの接続」をご参照ください。
-
次のステートメントを実行して、
test_adbpg_tableテーブルをクエリします。SELECT * FROM test_adbpg_table;データは期待どおりに AnalyticDB for PostgreSQL に書き込まれます。 次の図は、クエリ結果のサンプルを示しています。
クエリ結果には、
b1(int4)、b2(int4)、b3(text) の 3 つの列が含まれています。 複数行のデータが返されることから、データが宛先テーブルに同期されたことを確認できます。