このトピックでは、Realtime Compute for Apache Flink コンソールを使用して、Kafka から Hologres へリアルタイムのログデータをインジェストするデータ同期ジョブを手早く作成する方法を説明します。
前提条件
-
RAM ユーザーまたは RAM ロールに、Realtime Compute for Apache Flink コンソールにアクセスするために必要な権限が付与されていること。詳細については、「権限」をご参照ください。
-
Realtime Compute for Apache Flink ワークスペースが作成されていること。詳細については、「ワークスペースの作成」をご参照ください。
-
上流および下流ストレージ
-
ApsaraMQ for Kafka インスタンスが作成されていること。詳細については、「手順 2:インスタンスの購入とデプロイ」をご参照ください。
-
Hologres インスタンスが作成されていること。詳細については、「Hologres インスタンスの購入」をご参照ください。
説明ApsaraMQ for Kafka インスタンスと Hologres インスタンスは、Realtime Compute for Apache Flink ワークスペースと同じリージョンおよび VPC にある必要があります。そうでない場合は、ネットワーク接続を確立する必要があります。詳細については、「VPC をまたいだサービスへのアクセス」または「インターネットへのアクセス」をご参照ください。
-
ステップ 1:IP ホワイトリストの設定
Flink から Kafka と Hologres のインスタンスにアクセスできるように、Flink ワークスペースの CIDR ブロックを Kafka と Hologres の IP ホワイトリストに追加します。
-
Flink ワークスペースの VPC CIDR ブロックを取得します。
-
対象の[ワークスペース]の [Actions] 列で、 を選択します。
-
[Workspace Details] ダイアログボックスで、VSwitch の [CIDR block] を確認します。
ダイアログボックスには、ワークスペースの基本情報と VSwitch のリストが表示されます。[CIDR block] 列で、各アベイラビリティーゾーンの CIDR ブロックを確認します。次のステップのために、この CIDR ブロックを控えてください。
-
Flink ワークスペースの CIDR ブロックを Kafka インスタンスの IP ホワイトリストに追加します。
VPC エンドポイントの IP ホワイトリストを設定する必要があります。詳細な手順については、「IP ホワイトリストの設定」をご参照ください。ホワイトリスト編集ダイアログボックスで、[Add Whitelist IP] をクリックして CIDR ブロックを追加します。
-
Flink ワークスペースの CIDR ブロックを Hologres インスタンスの IP ホワイトリストに追加します。
Hologres インスタンスにログインし、その IP ホワイトリストを設定します。詳細な手順については、「IP ホワイトリスト」をご参照ください。HoloWeb の [Security Center] のホワイトリスト設定ページで、[Edit IP Whitelist] ダイアログボックスの [IP Address] フィールドに CIDR ブロックを入力し、[OK] をクリックします。
ステップ 2:Kafka テストデータの準備
Realtime Compute for Apache Flink の開発コンソールで、Faker コネクタを使用して ApsaraMQ for Kafka にデータを生成して書き込むには、次の手順に従います。
-
ApsaraMQ for Kafka コンソールで、「users」というトピックを作成します。
詳細については、「トピックの作成」をご参照ください。
-
ApsaraMQ for Kafka にデータを書き込むジョブを作成します。
-
対象のワークスペースの [Actions] 列で、[Console] をクリックします。
-
左側メニューで、 を選択します。
-
アイコンをクリックし、[New Blank Stream Draft] をクリックします。[Name] を入力し、[Engine Version] を選択します。Realtime Compute for Apache Flink には、さまざまなコードおよびデータ同期テンプレートも用意されています。各テンプレートには、特定のユースケース、コードサンプル、手順が含まれています。テンプレートをクリックすると、Realtime Compute for Apache Flink の機能と構文をすばやく習得し、ビジネスロジックを実装できます。詳細については、「コードテンプレート」および「データ同期テンプレート」をご参照ください。
パラメータ
説明
例
[Name]
ジョブの名前。
説明ジョブ名は、現在のワークスペース内で一意である必要があります。
flink-test
[Engine Version]
現在のジョブの Flink エンジンバージョン。
信頼性とパフォーマンスを向上させるには、[Recommended] または [Stable] のラベルが付いたバージョンを選択してください。エンジンバージョンの詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。
vvr-8.0.8-flink-1.17
-
[Create] をクリックします。
-
ジョブの SQL ステートメントを記述します。
次のコードをエディタにコピーし、環境に合わせてパラメータを変更してください。
CREATE TEMPORARY TABLE source ( id INT, first_name STRING, last_name STRING, `address` ROW<`country` STRING, `state` STRING, `city` STRING>, event_time TIMESTAMP ) WITH ( 'connector' = 'faker', 'number-of-rows' = '100', 'rows-per-second' = '10', 'fields.id.expression' = '#{number.numberBetween ''0'',''1000''}', 'fields.first_name.expression' = '#{name.firstName}', 'fields.last_name.expression' = '#{name.lastName}', 'fields.address.country.expression' = '#{address.country}', 'fields.address.state.expression' = '#{address.state}', 'fields.address.city.expression' = '#{address.city}', 'fields.event_time.expression' = '#{date.past ''15'',''SECONDS''}' ); CREATE TEMPORARY TABLE sink ( id INT, first_name STRING, last_name STRING, `address` ROW<`country` STRING, `state` STRING, `city` STRING>, `timestamp` TIMESTAMP METADATA ) WITH ( 'connector' = 'kafka', 'properties.bootstrap.servers' = 'alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092', 'topic' = 'users', 'format' = 'json', 'properties.enable.idempotence'='false' ); INSERT INTO sink SELECT id, first_name, last_name, `address` FROM source;次の表に、変更する必要があるパラメータを示します。
パラメータ
値の例
説明
properties.bootstrap.servers
alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092
ApsaraMQ for Kafka ブローカーエンドポイント。
カンマで区切られた、「host:port」形式のブローカーエンドポイントのリスト。VPC の [Domain Name] エンドポイントは、[Instance Details] ページの [Endpoint Information] セクションで確認できます。
topic
users
ApsaraMQ for Kafka トピック名。
-
ジョブを起動します。
-
ページで、[Deploy] をクリックします。
-
[Deploy draft] ダイアログボックスで、[Confirm] をクリックします。
-
ジョブのリソースを設定します。詳細については、「ジョブのリソースの設定」をご参照ください。
-
ページで、対象のデプロイメントを見つけ、[Actions] 列の [Start] をクリックします。起動設定の詳細については、「デプロイメントの起動」をご参照ください。
-
[Deployments] ページで、デプロイメントのランタイム情報とステータスを監視できます。
Faker ソースは有界ストリームを生成するため、デプロイメントのステータスは起動後約 1 分で FINISHED に変わります。デプロイメントが完了すると、データが「users」トピックに書き込まれています。以下は、JSON 形式のデータサンプルの例です。
{ "id": 765, "first_name": "Barry", "last_name": "Pollich", "address": { "country": "United Arab Emirates", "state": "Nevada", "city": "Powlowskifurt" } }
-
ステップ 3:データ同期ジョブの作成と開始
Flink CDC
-
Realtime Compute for Apache Flink 開発コンソールにログインして、データ同期ジョブを作成します。
-
Realtime Compute for Apache Flink コンソール にログインします。
-
対象のワークスペースの [Actions] 列で、[Console] をクリックします。
-
左側のナビゲーションペインで、 を選択します。
-
アイコンをクリックし、[新規下書き] をクリックします。[name] を入力し、[engine version] を選択します。パラメータ
説明
例
[名前]
ジョブの名前。
説明ジョブ名はプロジェクト内で一意である必要があります。
flink-test
[エンジンバージョン]
ジョブの Flink エンジンバージョン。
[Recommended] または [Stable] のラベルが付いたバージョンを選択します。これらのバージョンは、より高い信頼性とパフォーマンスを提供します。エンジンバージョンの詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。
vvr-8.0.8-flink-1.17
-
[Create] をクリックします。
-
-
Flink CDC ジョブを記述します。次のコードをエディターにコピーし、ご使用の環境に合わせてパラメータを更新します。
次のジョブは、Kafka の users トピックから JSON 形式のテーブルデータを、Hologres の flink_test_db データベースの test_schema スキーマにある users テーブルに同期します。
source: type: kafka name: Kafka Source properties.bootstrap.servers: alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092 topic: users scan.startup.mode: earliest-offset value.format: json json.infer-schema.flatten-nested-columns.enable: true sink: type: hologres name: Hologres Sink endpoint: hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80 dbname: flink_test_db username: ****** password: ** sink.type-normalize-strategy: ONLY_BIGINT_OR_TEXT transform: - source-table: \.*.\.* projection: \* primary-keys: id route: - source-table: users sink-table: test_schema.users次の表に、変更する必要があるパラメータを示します。
パラメータ
例
説明
properties.bootstrap.servers
alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092
Kafka ブローカーのアドレス。
形式は、host:port 形式のエントリをカンマで区切ったリストです。[インスタンス詳細] ページの [Network Information] セクションから、VPC ネットワークタイプの [Domain Name Endpoint] を取得できます。
topic
users
Kafka トピックの名前。
endpoint
hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80
Hologres インスタンスのエンドポイント。
形式は <ip>:<port> です。Hologres コンソールの [インスタンス詳細] ページの [Network Information] セクションから VPC エンドポイントを取得できます。
username
**
Hologres データベースのユーザー名とパスワード。Alibaba Cloud アカウントの AccessKey ID と AccessKey Secret を入力します。
重要AccessKey ペアの漏洩を防ぐために、変数管理を使用して AccessKey ID と AccessKey Secret を指定することを推奨します。詳細については、「変数の管理」をご参照ください。
password
**
dbname
flink_test_db
Hologres データベースの名前。
source-table
users
ソーステーブル。デフォルトでは、トピック名が使用されます。
sink-table
test_schema.users
宛先テーブル。テーブルを
schema.table_name形式で指定します。 -
[Save] をクリックします。
-
ページで、[Deploy] をクリックします。
-
ページで、対象のデプロイメントを見つけ、[Actions] 列の [Start] をクリックします。ジョブの起動設定の詳細については、「ジョブの開始」をご参照ください。
ジョブが開始されると、[Deployments] ページでそのランタイム情報とステータスを表示できます。このページには、[Status]、[ヘルススコア]、[CPU]、[memory] などのメトリクスを含むデプロイメントのリストが表示され、[Start] や [Stop] などのアクションを実行できます。
SQL
-
Realtime Compute for Apache Flink 開発コンソールにログインして、データ同期ジョブを作成します。
-
Realtime Compute for Apache Flink コンソール にログインします。
-
対象のワークスペースの [Actions] 列で、[Console] をクリックします。
-
左側のナビゲーションペインで、 を選択し、次に [新規] をクリックします。
-
アイコンをクリックし、次に [新規ブランクストリームドラフト] をクリックします。[名前] を入力し、[エンジンバージョン] を選択します。パラメータ
説明
例
[名前]
ジョブの名前。
説明ジョブ名はプロジェクト内で一意である必要があります。
flink-test
[エンジンバージョン]
ジョブの Flink エンジンバージョン。
[Recommended] または [Stable] のラベルが付いたバージョンを選択します。これらのバージョンは、より高い信頼性とパフォーマンスを提供します。エンジンバージョンの詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。
vvr-8.0.8-flink-1.17
-
[Create] をクリックします。
-
-
SQL ジョブを記述します。次のコードを SQL エディターにコピーし、ご使用の環境に合わせてパラメータを更新します。
INSERT INTO ステートメントを使用して、Kafka の users トピックから Hologres の flink_test_db データベースの users テーブルにデータを同期できます。
Hologres は、JSON および JSONB データ型に対して特別な最適化を提供します。INSERT INTO ステートメントを使用して、ネストされた JSON データを Hologres に書き込むことができます。
この方法では、まず Hologres に users テーブルを作成し、次に以下の SQL ステートメントを実行してテーブルにデータを書き込む必要があります。
CREATE TEMPORARY TABLE kafka_users ( `id` INT NOT NULL, `address` STRING, -- この列のデータはネストされた JSON です。 `offset` BIGINT NOT NULL METADATA, `partition` BIGINT NOT NULL METADATA, `timestamp` TIMESTAMP METADATA, `date` AS CAST(`timestamp` AS DATE), `country` AS JSON_VALUE(`address`, '$.country') ) WITH ( 'connector' = 'kafka', 'properties.bootstrap.servers' = 'alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092', 'topic' = 'users', 'format' = 'json', 'json.infer-schema.flatten-nested-columns.enable' = 'true', -- ネストされた列を自動的に展開します。 'scan.startup.mode' = 'earliest-offset' ); CREATE TEMPORARY TABLE holo ( `id` INT NOT NULL, `address` STRING, `offset` BIGINT, `partition` BIGINT, `timestamp` TIMESTAMP, `date` DATE, `country` STRING ) WITH ( 'connector' = 'hologres', 'endpoint' = 'hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80', 'username' = '******', 'password' = '******', 'dbname' = 'flink_test_db', 'tablename' = 'users' ); INSERT INTO holo SELECT * FROM kafka_users;次の表に、変更する必要があるパラメータを示します。
パラメータ
例
説明
properties.bootstrap.servers
alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092
Kafka ブローカーのアドレス。
形式は、host:port 形式のエントリをカンマで区切ったリストです。[インスタンス詳細] ページの [Network Information] セクションから、VPC ネットワークタイプの [Domain Name Endpoint] を取得できます。
topic
users
Kafka トピックの名前。
endpoint
hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80
Hologres インスタンスのエンドポイント。
形式は <ip>:<port> です。Hologres コンソールの [インスタンス詳細] ページの [Network Information] セクションから VPC エンドポイントを取得できます。
username
******
Hologres データベースのユーザー名とパスワード。Alibaba Cloud アカウントの AccessKey ID と AccessKey Secret を入力します。
重要AccessKey ペアの漏洩を防ぐために、変数管理を使用して AccessKey ID と AccessKey Secret を指定することを推奨します。詳細については、「変数の管理」をご参照ください。
password
******
dbname
flink_test_db
Hologres データベースの名前。
tablename
users
Hologres テーブルの名前。
説明-
INSERT INTO ステートメントを使用してデータを同期する場合、宛先データベースに users テーブルとその列を事前に作成する必要があります。
-
スキーマが public でない場合は、tablename パラメータを
schema.table_name形式で指定する必要があります。
-
-
[Save] をクリックします。
-
ページで、[Deploy] をクリックします。
-
ページで、対象のデプロイメントを見つけ、[Actions] 列の [Start] をクリックします。ジョブの起動設定の詳細については、「ジョブの開始」をご参照ください。
ジョブが開始されると、[Deployments] ページでそのランタイム情報とステータスを表示できます。このページには、[Status]、[ヘルススコア]、[CPU]、[memory] などのメトリクスを含むデプロイメントのリストが表示され、[Start] や [Stop] などのアクションを実行できます。
ステップ 4: 完全同期結果の表示
-
Hologres 管理コンソールにログインします。
-
[Instances] ページで、ターゲットインスタンスの名前をクリックします。
-
ページの右上隅で、[Connect to Instance] をクリックします。
-
[Metadata Management] タブで、flink_test_db データベース内の同期された users テーブルのテーブルスキーマとデータを表示します。
左側のナビゲーションウィンドウで、インスタンス名 > [flink_test_db] > [test_schema] > [Tables] を展開します。同期された [users] テーブルが表示されます。
同期されたテーブルのスキーマとデータは次のとおりです。
-
テーブルスキーマ
users テーブル名をダブルクリックして、テーブルスキーマを表示します。
users テーブルスキーマには、次のフィールドが含まれます: id (BIGINT、プライマリキー)、first_name (TEXT)、last_name (TEXT)、address.country (TEXT)、address.state (TEXT)、address.city (TEXT)。
説明完全データ同期中は、Kafka メタデータのパーティションとオフセットを Hologres テーブルのプライマリキーとして定義することを推奨します。これにより、ジョブがフェイルオーバーしてデータを再送信する際のデータ重複を防ぐことができます。
-
テーブルデータ
users テーブルページの右上隅で、[Query table] をクリックします。次のコマンドを入力し、[Run] をクリックします。
SELECT * FROM test_schema.users;コマンドは次の結果を返します:
クエリは複数の行を返します。これにより、レコードが users テーブルに正常に同期されたことが確認できます。返された各行には、id、first_name、last_name、address.country、address.state、address.city 列の完全なデータが含まれています。
-
ステップ 5:自動スキーマ同期の確認
-
ApsaraMQ for Kafka コンソールで、新しい列を含むメッセージを手動で送信します。
-
ApsaraMQ for Kafka コンソールにログインします。
-
[Instances] ページで、対象インスタンスの名前をクリックします。
-
[Topics] ページで、対象トピック (users) の名前をクリックします。
-
[Send Message] をクリックします。
-
メッセージを設定します。
[Start to Send and Consume Message] ダイアログボックスで、パラメーターを次のように設定します:
パラメーター
例
[送信方法]
[Console] を選択します。
[メッセージキー]
「flinktest」と入力します。
[メッセージ内容]
以下の JSON コンテンツをコピーして、「メッセージ内容」フィールドに貼り付けてください。
{ "id": 100001, "first_name": "Ichiro", "last_name": "Tanaka", "address": { "country": "Isle of Man", "state": "Montana", "city": "East Coleburgh" }, "house-points": { "house": "Pukwudgie", "points": 76 } }説明この例では、 house-points は新しいネストされた列です。
[指定パーティションへの送信]
[Yes] を選択します。
[パーティション ID]
「0」と入力します。
-
[OK] をクリックします。
-
-
Hologres コンソールで、users テーブルのスキーマとデータの変更を確認します。
-
Hologres コンソールにログインします。
-
[Instances] ページで、対象インスタンスの名前をクリックします。
-
ページの右上隅にある [インスタンスに接続] をクリックします。
-
[メタデータ管理] タブで、users テーブル名をダブルクリックします。
-
[テーブルのクエリ] をクリックし、次の文を入力して [Run] をクリックします。
SELECT * FROM test_schema.users; -
クエリ結果を表示します。
クエリは次の結果を返します:
この結果から、ID 100001 のレコードが Hologres に正常に書き込まれ、2 つの新しい列 (house-points.house と house-points.points) が Hologres テーブルに追加されたことがわかります。
説明ApsaraMQ for Kafka に送信されたメッセージには、ネストされた列が 1 つ (house-points) しか含まれていません。しかし、WITH 句で json.infer-schema.flatten-nested-columns.enable が指定されているため、Realtime Compute for Apache Flink はこの列を自動的にフラット化し、ネストされたフィールドアクセスパスを新しい列名として使用します。
-
参考資料
-
Message Queue for Apache Kafka をソーステーブルまたは結果テーブルとして使用する方法については、「Message Queue for Apache Kafka」をご参照ください。
-
ノードの並列度とリソースを調整してジョブのパフォーマンスを向上させるには、「ジョブデプロイメントの設定」をご参照ください。