Realtime Compute for Apache Flink は、フル同期から増分同期への自動切り替え、メタデータ検出、スキーマエボリューション、データベース全体の同期を自動的に処理することで、リアルタイムのデータインジェストを簡素化します。このトピックでは、ApsaraDB RDS for MySQL から Hologres へデータをストリーミングするデータインジェストジョブを迅速に構築する方法を説明します。
背景情報
DMS コンソールでは、これらのデータベースとテーブルは以下のように表示されます。
RDS インスタンスには、複数のデータベース (tpc_ds、tpc_ds_large、user_db1、user_db2、user_db3、__recycle_bin__) が含まれています。user_db3 データベースには、user03、user06、user09 の 3 つのテーブルが含まれています。user03 テーブルには、id (int(11)) と name (varchar(255)) の 2 つの列があります。
以下の手順に従って、これらすべてのテーブルを Hologres に同期し、シャード化されたユーザーテーブルを単一のテーブルにマージするデータインジェストジョブを開発します。
このトピックでは、Flink CDC データインジェストを使用して、データベース全体の同期とシャード化されたテーブルのマージを実行します。これにより、単一のジョブで、フル同期と増分同期、およびリアルタイムのスキーマ変更同期を完了できます。
前提条件
-
RAM ユーザーまたは RAM ロールを使用する場合、Flink コンソールにアクセスするために必要な権限があることを確認してください。詳細については、「権限の管理」をご参照ください。
-
Flink ワークスペースを作成済みであること。詳細については、「Realtime Compute for Apache Flinkの有効化」をご参照ください。
-
アップストリームおよびダウンストリームストレージ
-
ApsaraDB RDS for MySQL インスタンスを作成済みであること。詳細については、「(非推奨、"手順1" にリダイレクト) ApsaraDB RDS for MySQLインスタンスをすばやく作成する」をご参照ください。
-
Hologres インスタンスを作成済みであること。詳細については、「Hologresインスタンスの購入」をご参照ください。
説明ApsaraDB RDS for MySQL インスタンスと Hologres インスタンスは、Flink ワークスペースと同じリージョンおよび Virtual Private Cloud (VPC) 内にある必要があります。そうでない場合は、ネットワーク接続を確立する必要があります。詳細については、「VPCをまたいで他のサービスにアクセスするにはどうすればよいですか?」および「インターネットにアクセスするにはどうすればよいですか?」をご参照ください。
-
-
テストデータを準備し、IP ホワイトリストを設定済みであること。詳細については、「MySQLテストデータとHologresデータベースの準備」および「IP ホワイトリストの設定」をご参照ください。
MySQL テストデータと Hologres データベースの準備
-
tpc_ds.sql、user_db1.sql、user_db2.sql、user_db3.sql をクリックして、テストデータファイルをローカルマシンにダウンロードします。
-
DMS コンソールで、ApsaraDB RDS for MySQL インスタンスにテストデータを準備します。
-
DMS を使用して ApsaraDB RDS for MySQL インスタンスにログインします。
詳細については、「(非推奨、"手順2" にリダイレクト) DMSを使用してApsaraDB RDS for MySQLインスタンスにログインする」をご参照ください。
-
SQL コンソール ウィンドウに、次のコマンドを入力し、[実行] をクリックします。
次のコマンドは、
tpc_ds、user_db1、user_db2、user_db3の 4 つのデータベースを作成します。CREATE DATABASE tpc_ds; CREATE DATABASE user_db1; CREATE DATABASE user_db2; CREATE DATABASE user_db3; -
トップナビゲーションバーで、[データインポート] をクリックします。
-
[一括データインポート] タブで、データをインポートするデータベースを選択し、対応する SQL ファイルをアップロードして [送信] をクリックした後、[変更を実行] をクリックします。表示されたダイアログボックスで、[実行を確認] をクリックします。
この操作を繰り返して、対応するデータファイルを
tpc_ds、user_db1、user_db2、user_db3データベースにインポートします。[ファイルエンコーディング] を [自動検出] に、[インポートモード] を [高速モード] または [セーフモード] に、[ファイルタイプ] を [SQL スクリプト]、[CSV]、または [Excel] に設定します。添付ファイルは txt/sql/csv/xlsx/zip 形式に対応しており、サイズは最大 5 GB です。
-
-
Hologres コンソールで、マージされたユーザーテーブルデータを格納するために、
my_userという名前のデータベースを作成します。詳細については、「データベースの作成」をご参照ください。
IP ホワイトリストの設定
Flink が ApsaraDB RDS for MySQL および Hologres インスタンスにアクセスできるようにするには、Flink ワークスペースの CIDR ブロックを両方のインスタンスの IP ホワイトリストに追加します。
-
Flink ワークスペースの CIDR ブロックを取得します。
-
ワークスペースのリストで、目的の [ワークスペース] を見つけ、[アクション] 列で を選択します。
-
[ワークスペース詳細] ダイアログボックスで、 Flink [vSwitch] の [CIDR ブロック]情報を確認します。
-
Flink の CIDR ブロックを ApsaraDB RDS for MySQL インスタンスの IP ホワイトリストに追加します。
詳細については、「IP ホワイトリストの設定」をご参照ください。
[ホワイトリストグループの変更] ダイアログボックスで、[グループホワイトリスト] フィールドに Flink フルマネージドの CIDR ブロックを入力し、[OK] をクリックします。
-
Flink の CIDR ブロックを Hologres インスタンスの IP ホワイトリストに追加します。
HoloWeb でデータ接続を設定する際、接続の IP ホワイトリストを設定する前に、[Login Method] を [Password-free login for current user] に設定します。詳細については、「IP ホワイトリスト」をご参照ください。
HoloWeb で、上部のナビゲーションバーから [セキュリティセンター] を選択し、左側のペインで [IP ホワイトリスト] をクリックし、[IP ホワイトリストの編集] ダイアログボックスで次のパラメーターを設定します。
-
[グループ]: デフォルトを選択します
-
[データベース制限]:すべてを選択
-
[ユーザー制限]:すべて選択
-
[IP アドレス]:Flink CIDR ブロックを CIDR 表記で入力します(例: 172.xx.0/19)。 複数の IP アドレスは改行で区切ります。
適用するには、[OK] をクリックします。
-
手順1:データインジェストジョブの開発
-
Flink 開発コンソールにログインし、新しいジョブを作成します。
-
ページで、[新規作成] をクリックします。
-
[空白のデータ取り込みドラフト] をクリックします。
Flink は、豊富なコードテンプレートセットを提供しており、それぞれに特定のユースケース、コードサンプル、ガイダンスが含まれています。テンプレートをクリックすると、ビジネスロジックを実装するための製品機能や構文について学習できます。
-
[次へ] をクリックします。
-
[新規ドラフト] ダイアログボックスで、パラメーターを設定します。
パラメーター
説明
例
[名前]
ジョブの名前。
説明ジョブ名は、現在のプロジェクト内で一意である必要があります。
flink-test
[場所]
ジョブのコードファイルが保存されるフォルダー。
既存のフォルダーの横にある
アイコンをクリックして、サブフォルダーを作成することもできます。Job Drafts
[エンジンバージョン]
ジョブで使用される Flink エンジンのバージョン。エンジンのバージョン番号、バージョンの互換性、および重要なライフサイクル日については、「エンジンバージョン」をご参照ください。
vvr-11.1-jdk11-flink-1.20
-
[OK] をクリックします。
-
-
次のジョブコードをジョブエディターにコピーします。
次のコードは、
tpc_dsデータベースのすべてのテーブルを Hologres に同期し、シャード化されたユーザーテーブルを Hologres の単一のテーブルにマージします。source: type: mysql name: MySQL Source hostname: localhost port: 3306 username: username password: password tables: tpc_ds.\.*,user_db[0-9]+.user[0-9]+ server-id: 8601-8604 # (任意) テーブルと列のコメントを同期します。 include-comments.enabled: true # (任意) TaskManager の OutOfMemory エラーの可能性を防ぐために、無制限チャンクの配布を優先します。 scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (任意) 読み取りを高速化するために、解析フィルターを有効にします。 scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: hologres name: Hologres Sink endpoint: ****.hologres.aliyuncs.com:80 dbname: cdcyaml_test username: ${secret_values.holo-username} password: ${secret_values.holo-password} sink.type-normalize-strategy: BROADEN route: # シャード化されたユーザーテーブルを my_user.users テーブルにマージして同期します。 - source-table: user_db[0-9]+.user[0-9]+ sink-table: my_user.users説明MySQL の
tpc_dsデータベースのテーブルは、同期先で同じ名前のテーブルに直接マッピングされるため、routeセクションに追加のマッピング設定は必要ありません。ods_tps_dsのような別の名前のデータベースにテーブルを同期するには、routeモジュールを次のように設定します。route: # シャード化されたユーザーテーブルを my_user.users テーブルにマージして同期します。 - source-table: user_db[0-9]+.user[0-9]+ sink-table: my_user.users # tpc_ds 配下のすべてのテーブルのデータベース名を変更し、ods_tps_ds に同期します。 - source-table: tpc_ds.\.* sink-table: ods_tps_ds.<> replace-symbol: <>
手順2:ジョブの開始
-
ページで、[デプロイ] をクリックします。表示されるダイアログボックスで、[確認] をクリックします。
ダイアログボックスで、[備考] を入力し、[ジョブタグ] を設定して、[デプロイターゲット] ドロップダウンリストからキューを選択できます (デフォルト:
default-queue)。 注:デプロイは、次回のジョブ起動時に有効になります。 -
ページで、対象のジョブの [アクション] 列にある [開始] をクリックします。必要に応じてパラメーターを設定します。詳細については、「ジョブを開始する」をご参照ください。
-
[開始] をクリックします。
ジョブが開始されると、「ジョブオペレーション」ページでそのランタイム情報とステータスを表示できます。ジョブのステータスには、[実行中]、[失敗]、[停止] があります。上部の [ストリームジョブ] ドロップダウンフィルターを使用して、ジョブリストをタイプ別にフィルターできます。
手順3:フル同期結果の検証
-
Hologres 管理コンソールにログインします。
-
[メタデータ管理] タブで、Hologres インスタンスの
tpc_dsデータベースに 24 個のテーブルとそのデータが存在することを確認します。インスタンス > [tpc_ds] > [public] > [テーブル] に移動すると、call_center、catalog_sales、customer、store_sales などのテーブルが表示されます。 [store_sales] テーブルを選択し、[データプレビュー] タブをクリックすると、ss_sold_date、ss_sold_time、ss_item_sk、ss_customer などの列が表示されます。
-
[メタデータ管理] タブで、
my_userデータベースのusersテーブルのスキーマを確認します。同期されたテーブルのスキーマとデータは以下のとおりです。
-
テーブルスキーマ
同期後、
usersテーブルには次の列が含まれます。-
_db_name(text, プライマリキー) -
_table_name(text, プライマリキー) -
id(int4, プライマリキー) -
name(varchar 255, null 許容)
_db_name、_table_name、idの列が複合プライマリキーを構成します。usersテーブルのスキーマには、ソースの MySQL テーブルにはない 2 つの追加列、_db_nameと_table_nameが含まれています。これらの列は各行のソースデータベースとソーステーブルを示し、複合プライマリキーの一部として、シャード化されたテーブルがマージされた後のデータの一意性を保証します。 -
-
テーブルデータ
usersテーブル情報ページの右上隅で、[テーブルのクエリ] をクリックします。次のコマンドを入力し、[実行] をクリックします。select * from users order by _db_name,_table_name,id;クエリ結果から、ユーザーデータが 3 つのデータベース (
user_db1(user01、user04、user07)、user_db2(user02、user05、user08)、user_db3(user03、user06、user09)) に完全に同期され、_db_name、_table_name、id、nameの列を持つ合計 9 件のレコードがあることがわかります。
-
手順4:増分同期結果の検証
フル同期が完了すると、ジョブは手動介入なしで自動的に増分同期フェーズに切り替わります。[Monitoring and Alerts] タブで currentEmitEventTimeLag の値を確認することで、データ同期フェーズを判断できます。
-
対象のワークスペースの [アクション] 列で、[コンソール] をクリックします。
-
ページで、対象のジョブの名前をクリックします。
-
[アラーム] (または [メトリクス]) タブをクリックします。
-
currentEmitEventTimeLagチャートを調べて、データ同期フェーズを判断します。
-
値が 0 の場合は、フル同期フェーズを示します。
-
値が 0 より大きい場合は、増分同期フェーズを示します。
-
-
リアルタイムのデータとスキーマ変更の同期を検証します。
MySQL CDC ソースは、増分フェーズ中のリアルタイムのデータとスキーマの同期をサポートしています。これを検証するには、ジョブがこのフェーズに入った後、MySQL でシャード化されたユーザーテーブルのスキーマとデータを変更します。
-
DMS を使用して ApsaraDB RDS for MySQL インスタンスにログインします。
詳細については、「(非推奨、"手順2" にリダイレクト) DMSを使用してApsaraDB RDS for MySQLインスタンスにログインする」をご参照ください。
-
user_db2データベースで、次のコマンドを実行してuser02テーブルのスキーマを変更し、データを挿入および更新します。USE `user_db2`; ALTER TABLE `user02` ADD COLUMN `age` INT; -- age 列を追加します。 INSERT INTO `user02` (id, name, age) VALUES (27, 'Tony', 30); -- age データを含む行を挿入します。 UPDATE `user05` SET name='JARK' WHERE id=15; -- 別のテーブルを更新し、名前を大文字に変更します。 -
Hologres コンソールで、
usersテーブルのスキーマとデータの変更を確認します。usersテーブル情報ページの右上隅で、[テーブルをクエリ] をクリックして次のコマンドを入力し、[実行] をクリックします。select * from users order by _db_name,_table_name,id;シャード化されたテーブルのスキーマが異なっていても、
user02のスキーマ変更とデータ変更はリアルタイムで伝播されます。Hologres のusersテーブルには、新しいage列、挿入された Tony のレコード、および更新された JARK のレコードが表示されます。クエリ結果には、user_db1、user_db2、user_db3 にまたがる users テーブルから 10 行が表示されます。増分同期の変更が反映されています。user_db2 のユーザー Tony (id 27) は
ageが 30 に更新され、別のレコードはnameが JARK に変更されています。残りの行の age の値は \N であり、増分同期が機能していることを確認できます。
-
(任意) 手順5:ジョブリソースの設定
パフォーマンスを向上させるために、データ量に基づいて、並列度、TaskManager のメモリ、CU などのジョブリソースを調整できます。
-
ページで、対象のジョブの名前をクリックします。
-
[設定] タブで、[リソース] セクションの右上隅にある [編集] をクリックします。
-
TaskManager のメモリや並列度などのリソースパラメーターを手動で設定します。
-
[リソース] セクションの右側で、[保存] をクリックします。
-
ジョブを再起動します。
リソース設定の変更は、ジョブを再起動した後にのみ有効になります。
関連ドキュメント
-
各データインジェストモジュールの構文については、「Flink CDCデータインジェストジョブ開発リファレンス」をご参照ください。
-
データインジェストジョブの実行中に問題が発生した場合は、「データインジェストジョブの一般的な問題と解決策」をご参照ください。