ユーザー行動データの処理は、その膨大な量と多様なフォーマットのため、困難を伴います。従来のワイドテーブルモデルは、効率的なクエリパフォーマンスを提供する一方で、高いデータ冗長性、ストレージオーバーヘッドの増加、メンテナンスサイクルの長期化といった代償を伴います。このチュートリアルでは、Realtime Compute for Apache Flink、ApsaraDB for MongoDB、Hologres を使用して、これらのトレードオフに対応するリアルタイムのワイドテーブルパイプラインを構築する方法を説明します。
仕組み
Realtime Compute for Apache Flink はストリーム処理を担います。ApsaraDB for MongoDB は、柔軟なスキーマと高い読み書きスループットを持つドキュメント指向の NoSQL データベースとして、ファクトテーブルとディメンションテーブルを格納します。Hologres は分析用データウェアハウスとして機能し、データは書き込み後すぐにクエリ可能になります。
このパイプラインは、Kafka トピックで接続された 2 つの Flink ジョブを使用します。
-
ジョブ 1 は MongoDB からの変更データキャプチャ (CDC) ストリームを読み取ります。ファクトテーブル (
game_sales) が変更されると、そのプライマリキー (PK) は直接 Kafka にストリーミングされます。ディメンションテーブルが変更されると、ルックアップ結合によって影響を受けるファクトテーブルのレコードが特定され、その PK が Kafka に送信されます。 -
ジョブ 2 は Kafka から PK を消費し、MongoDB のファクトテーブルとディメンションテーブルに対してルックアップ結合を実行して完全なワイドテーブルレコードを再構築し、その結果を Hologres にアップサートします。
メリット:
-
高い書き込みスループット: ApsaraDB for MongoDB は、シャードクラスターでの高同時実行数の読み書きを処理し、パフォーマンスとストレージをスケーリングして、大量の頻繁な更新に対応します。
-
効率的な変更の伝播: 影響を受けるレコードの PK のみが Kafka に転送され、行全体は転送されません。これにより、総データ量に関係なく、再処理が最小限に抑えられます。
-
リアルタイムクエリ: Hologres は低レイテンシーのアップサートをサポートし、各書き込みの直後にデータをクエリ可能にします。
ハンズオン
このチュートリアルを完了すると、ファクトテーブルとディメンションテーブルの両方に対する MongoDB の変更が、Hologres のワイドテーブルに自動的に伝播し、即座にクエリ可能になる稼働中のパイプラインが完成します。
このパイプラインは、3 つの MongoDB コレクションを 1 つの Hologres ワイドテーブルに結合します。
game_sales
<table> <thead> <tr> <td><p>sale_id</p></td> <td><p>game_id</p></td> <td><p>platform_id</p></td> <td><p>sale_date</p></td> <td><p>units_sold</p></td> <td><p>sale_amt</p></td> <td><p>status</p></td> </tr> </thead> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <tbody></tbody> </table>
game_dimension
<table> <thead> <tr> <td><p>game_id</p></td> <td><p>game_name</p></td> <td><p>release_date</p></td> <td><p>developer</p></td> <td><p>publisher</p></td> </tr> </thead> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <tbody></tbody> </table>
platform_dimension
<table> <thead> <tr> <td><p>platform_id</p></td> <td><p>platform_name</p></td> <td><p>type</p></td> </tr> </thead> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <tbody></tbody> </table>
game_sales_details
<table> <thead> <tr> <td><p>sale_id</p></td> <td><p>game_id</p></td> <td><p>platform_id</p></td> <td><p>sale_date</p></td> <td><p>units_sold</p></td> <td><p>sale_amt</p></td> <td><p>status</p></td> </tr> </thead> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <colgroup></colgroup> <tbody> <tr> <td><p>game_name</p></td> <td><p>release_date</p></td> <td><p>developer</p></td> <td><p>publisher</p></td> <td><p>platform_name</p></td> <td><p>type</p></td> <td><p></p></td> </tr> </tbody> </table>
変更の伝播フロー
各変更は 4 つのステージを経て流れます。
-
キャプチャ: MongoDB ディメンションテーブルからのリアルタイムの変更を検出します。
-
伝播: ディメンションテーブルが変更されると、ジョブ 1 はルックアップ結合 (例えば、
game_idで) を実行してファクトテーブル内の影響を受ける行を見つけ、その PK (例えば、sale_id) を抽出します。 -
トリガー: PK を Kafka に送信して、保留中のリフレッシュをジョブ 2 に通知します。
-
アップサート: ジョブ 2 は最新のデータをフェッチし、ワイドテーブルの行を再構築して、Hologres にアップサートします。
前提条件
開始する前に、以下が揃っていることを確認してください。
-
VVR 8.0.5 以降を実行している Realtime Compute for Apache Flink ワークスペース。詳細については、「Realtime Compute for Apache Flink の有効化」をご参照ください。
-
バージョン 4.0 以降を実行している ApsaraDB for MongoDB インスタンス。詳細については、「シャードクラスターインスタンスの作成」をご参照ください。
-
バージョン 1.3 以降を実行している Hologres 専用インスタンス。詳細については、「Hologres インスタンスの購入」をご参照ください。
-
ApsaraMQ for Kafka インスタンス。詳細については、「ApsaraMQ for Kafka インスタンスのデプロイ」をご参照ください。
-
4 つのインスタンスすべてが同じ Virtual Private Cloud (VPC) 内にあること。異なる VPC にある場合は、VPC 間の接続を確立するか、Realtime Compute for Apache Flink のインターネットアクセスを有効にする必要があります。詳細については、「Realtime Compute for Apache Flink はどのようにして VPC をまたいでサービスにアクセスしますか?」および「Realtime Compute for Apache Flink はどのようにしてインターネットにアクセスしますか?」をご参照ください。
-
関連リソースに対する RAM ユーザーまたは RAM ロールの権限。
ステップ 1:データの準備
MongoDB コレクションの作成
-
ご利用の Flink ワークスペースの CIDR ブロックを MongoDB のホワイトリストに追加します。詳細については、「インスタンスのホワイトリストの設定」および「ホワイトリストの設定方法」をご参照ください。
-
Data Management (DMS) コンソールの SQL エディターで、
mongo_testデータベースを作成します。use mongo_test; -
game_sales、game_dimension、platform_dimensionコレクションを作成し、サンプルデータを挿入します。// ゲーム販売テーブル (ステータス:1 = アクティブ、0 = 論理削除) db.game_sales.insert( [ {sale_id:0,game_id:101,platform_id:1,"sale_date":"2024-01-01",units_sold:500,sale_amt:2500,status:1}, ] ); // ゲームディメンションテーブル db.game_dimension.insert( [ {game_id:101,"game_name":"SpaceInvaders","release_date":"2023-06-15","developer":"DevCorp","publisher":"PubInc"}, {game_id:102,"game_name":"PuzzleQuest","release_date":"2023-07-20","developer":"PuzzleDev","publisher":"QuestPub"}, {game_id:103,"game_name":"RacingFever","release_date":"2023-08-10","developer":"SpeedCo","publisher":"RaceLtd"}, {game_id:104,"game_name":"AdventureLand","release_date":"2023-09-05","developer":"Adventure","publisher":"LandCo"}, ] ); // プラットフォームディメンションテーブル db.platform_dimension.insert( [ {platform_id:1,"platform_name":"PCGaming","type":"PC"}, {platform_id:2,"platform_name":"PlayStation","type":"Console"}, {platform_id:3,"platform_name":"Mobile","type":"Mobile"} ] ); -
挿入を検証します。
db.game_sales.find(); db.game_dimension.find(); db.platform_dimension.find();db.game_dimension.find()クエリは 4 つのレコードを返します。-
game_id: 101,game_name: SpaceInvaders,release_date: 2023-06-15 -
game_id: 102,game_name: PuzzleQuest,release_date: 2023-07-20 -
game_id: 103,game_name: RacingFever,release_date: 2023-08-10 -
game_id: 104,game_name: AdventureLand,release_date: 2023-09-05
-
Hologres テーブルの作成
-
Hologres コンソールにログインし、左側のナビゲーションウィンドウで [インスタンス] をクリックし、ご利用の Hologres インスタンスをクリックします。右上隅の [インスタンスに接続] をクリックします。
-
上部のナビゲーションバーで、[メタデータ管理] > [データベースの作成] をクリックします。[データベース名] フィールドに
testと入力し、[ポリシー] を [SPM] に設定し、[OK] をクリックします。詳細については、「データベースの作成」をご参照ください。ダイアログボックスで、インスタンス [User-behavior-test] を選択し、[すぐにログオン] を [はい] に設定し、[OK] をクリックします。
-
上部のナビゲーションバーで、[SQL エディター] をクリックします。SQL アイコンをクリックして新しい SQL クエリを作成し、ターゲットインスタンスとデータベースを選択して、次の文を実行して
game_sales_detailsワイドテーブルを作成します。CREATE TABLE game_sales_details( sale_id INT not null primary key, game_id INT, platform_id INT, sale_date VARCHAR(50), units_sold INT, sale_amt INT, status INT, game_name VARCHAR(50), release_date VARCHAR(50), developer VARCHAR(50), publisher VARCHAR(50), platform_name VARCHAR(50), type VARCHAR(50) );
Kafka トピックの作成
-
ApsaraMQ for Kafka コンソールにログインします。左側のナビゲーションウィンドウで [インスタンス] をクリックし、ご利用のインスタンスをクリックします。
-
左側のナビゲーションウィンドウで、[ホワイトリスト管理] をクリックし、ご利用の Flink ワークスペースの CIDR ブロックを追加します。
-
左側のナビゲーションウィンドウで、[トピック] > [トピックの作成] をクリックします。右側のペインで、[名前] フィールドに
game_sales_factと入力し、説明を入力して、他のすべてのフィールドはデフォルト値のままにします。[OK] をクリックします。
ステップ 2:ストリームジョブの作成
ジョブ 1:プライマリキーの Kafka への書き込み
ジョブ 1 は 3 つすべての MongoDB コレクションをモニターします。game_sales が変更されると、sale_id は直接 Kafka にストリーミングされます。ディメンションテーブルが変更されると、game_sales に対するルックアップ結合によって影響を受ける sale_id の値が取得され、Kafka にストリーミングされます。
MongoDB コネクタはこのパイプラインで 2 つの役割を果たします。ジョブ 1 では、MongoDB の変更ストリームを読み取る CDC ソースとして機能します。ジョブ 2 では、プライマリキーによって各ドキュメントの現在の状態をフェッチするルックアップソースとして機能します。両方の役割で同じコネクタ構成が使用されます。
-
ご利用のワークスペースの [操作] 列で、[コンソール] をクリックします。
-
左側のナビゲーションメニューで、[開発] > [ETL] をクリックします。
-
[新しい空白のストリームドラフト] をクリックします。
-
[新しいドラフト] ダイアログで、[名前] に
dwd_mongo_kafkaと入力し、エンジンバージョンを選択して、[作成] をクリックします。 -
次の SQL をエディターにコピーします。3 つの
INSERT文はそれぞれ、独立して 1 つの MongoDB コレクションからの変更をキャプチャし、影響を受けるsale_idの値を Kafka sink にストリーミングします。これにより、いずれかのテーブルが変更されたときに、Hologres が正確かつリアルタイムに更新されることが保証されます。接続文字列やパスワードなどの機密性の高い値は、ハードコーディングするのではなく、変数として格納してください。詳細については、「変数の管理」をご参照ください。時間 game_id = 101のgame_dimensionT1 game_name = "SpaceInvaders"T2 game_name = "SpaceInvaders_v2"-- ソース:game_sales (CDC ストリーム) CREATE TEMPORARY TABLE game_sales ( `_id` STRING, -- MongoDB 自動生成 ID sale_id INT, -- 販売 ID PRIMARY KEY (_id) NOT ENFORCED ) WITH ( 'connector' = 'mongodb', 'uri' = '${secret_values.MongoDB-URI}', 'database' = 'mongo_test', 'collection' = 'game_sales' ); -- ソース:game_dimension (CDC ストリーム) CREATE TEMPORARY TABLE game_dimension ( `_id` STRING, game_id INT, PRIMARY KEY (_id) NOT ENFORCED ) WITH ( 'connector' = 'mongodb', 'uri' = '${secret_values.MongoDB-URI}', 'database' = 'mongo_test', 'collection' = 'game_dimension' ); -- ソース:platform_dimension (CDC ストリーム) CREATE TEMPORARY TABLE platform_dimension ( `_id` STRING, platform_id INT, PRIMARY KEY (_id) NOT ENFORCED ) WITH ( 'connector' = 'mongodb', 'uri' = '${secret_values.MongoDB-URI}', 'database' = 'mongo_test', 'collection' = 'platform_dimension' ); -- ルックアップソース:game_sales (ディメンション変更の結合に使用) CREATE TEMPORARY TABLE game_sales_dim ( `_id` STRING, sale_id INT, game_id INT, platform_id INT, PRIMARY KEY (_id) NOT ENFORCED ) WITH ( 'connector' = 'mongodb', 'uri' = '${secret_values.MongoDB-URI}', 'database' = 'mongo_test', 'collection' = 'game_sales' ); -- Sink:影響を受ける PK を格納する Kafka トピック CREATE TEMPORARY TABLE game_sales_fact ( sale_id INT, PRIMARY KEY (sale_id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'properties.bootstrap.servers' = '${secret_values.Kafka-hosts}', 'topic' = 'game_sales_fact', 'key.format' = 'json', 'value.format' = 'json', 'properties.enable.idempotence' = 'false' -- ApsaraMQ for Kafka への書き込み時に必須 ); BEGIN STATEMENT SET; -- game_sales の変更から PK をストリーミング INSERT INTO game_sales_fact (sale_id) SELECT sale_id FROM game_sales; -- game_dimension の変更によって影響を受ける game_sales 行の PK をストリーミング INSERT INTO game_sales_fact (sale_id) SELECT gs.sale_id FROM game_dimension AS gd JOIN game_sales_dim FOR SYSTEM_TIME AS OF PROCTIME() AS gs ON gd.game_id = gs.game_id; -- platform_dimension の変更によって影響を受ける game_sales 行の PK をストリーミング INSERT INTO game_sales_fact (sale_id) SELECT gs.sale_id FROM platform_dimension AS pd JOIN game_sales_dim FOR SYSTEM_TIME AS OF PROCTIME() AS gs ON pd.platform_id = gs.platform_id; END;ルックアップ結合について
FOR SYSTEM_TIME AS OF PROCTIME()句は、ルックアップ結合 (テンポラル結合の一種) を定義します。ソース行が処理される瞬間に、結合は MongoDB から一致するディメンションテーブルの行をフェッチし、そのスナップショットを結合結果としてフリーズします。後でディメンションテーブルが更新されても、すでに処理された行は影響を受けません。例えば、T1 で処理されたgame_sales行は"SpaceInvaders"と結合します。T2 で処理された行は"SpaceInvaders_v2"と結合します。結合条件はgd.game_id = gs.game_idおよびpd.platform_id = gs.platform_idです。詳細については、「ディメンションテーブルの JOIN 文」および「Kafka、Upsert Kafka、または Kafka JSON カタログの選択」をご参照ください。 -
右上隅で [デプロイ] をクリックします。ダイアログで [確認] をクリックします。詳細については、「ジョブのデプロイ」をご参照ください。
ジョブ 2:ワイドテーブルの再構築と Hologres へのアップサート
ジョブ 2 は game_sales_fact Kafka トピックから sale_id の値を消費し、MongoDB のファクトテーブルとディメンションテーブルに対してルックアップ結合を実行し、結果のワイドテーブル行を Hologres にアップサートします。
ジョブ 1 の手順に従って、dws_kafka_mongo_holo という名前の新しいドラフトを作成し、次の SQL でデプロイします。
-- ソース:影響を受ける PK を提供する Kafka トピック
CREATE TEMPORARY TABLE game_sales_fact
(
sale_id INT,
PRIMARY KEY (sale_id) NOT ENFORCED
)
WITH (
'connector' = 'upsert-kafka',
'properties.bootstrap.servers' = '${secret_values.Kafka-hosts}',
'topic' = 'game_sales_fact',
'key.format' = 'json',
'value.format' = 'json',
'properties.group.id' = 'game_sales_fact',
'properties.auto.offset.reset' = 'earliest'
);
-- ルックアップソース:game_sales ファクトテーブル
CREATE TEMPORARY TABLE game_sales
(
`_id` STRING,
sale_id INT,
game_id INT,
platform_id INT,
sale_date STRING,
units_sold INT,
sale_amt INT,
status INT,
PRIMARY KEY (_id) NOT ENFORCED
)
WITH (
'connector' = 'mongodb',
'uri' = '${secret_values.MongoDB-URI}',
'database' = 'mongo_test',
'collection' = 'game_sales'
);
-- ルックアップソース:game_dimension
CREATE TEMPORARY TABLE game_dimension
(
`_id` STRING,
game_id INT,
game_name STRING,
release_date STRING,
developer STRING,
publisher STRING,
PRIMARY KEY (_id) NOT ENFORCED
)
WITH (
'connector' = 'mongodb',
'uri' = '${secret_values.MongoDB-URI}',
'database' = 'mongo_test',
'collection' = 'game_dimension'
);
-- ルックアップソース:platform_dimension
CREATE TEMPORARY TABLE platform_dimension
(
`_id` STRING,
platform_id INT,
platform_name STRING,
type STRING,
PRIMARY KEY (_id) NOT ENFORCED
)
WITH (
'connector' = 'mongodb',
'uri' = '${secret_values.MongoDB-URI}',
'database' = 'mongo_test',
'collection' = 'platform_dimension'
);
-- Sink:Hologres ワイドテーブル
CREATE TEMPORARY TABLE IF NOT EXISTS game_sales_details
(
sale_id INT,
game_id INT,
platform_id INT,
sale_date STRING,
units_sold INT,
sale_amt INT,
status INT,
game_name STRING,
release_date STRING,
developer STRING,
publisher STRING,
platform_name STRING,
type STRING,
PRIMARY KEY (sale_id) NOT ENFORCED
)
WITH (
'connector' = 'hologres',
'dbname' = 'test',
'tablename' = 'public.game_sales_details',
'username' = '${secret_values.AccessKeyID}',
'password' = '${secret_values.AccessKeySecret}',
'endpoint' = '${secret_values.Hologres-endpoint}',
'sink.delete-strategy' = 'IGNORE_DELETE', -- 挿入または更新のみ。行は削除しない
'sink.on-conflict-action' = 'INSERT_OR_UPDATE', -- 部分的な列の更新を有効化
'sink.partial-insert.enabled' = 'true'
);
INSERT INTO game_sales_details (
sale_id, game_id, platform_id, sale_date, units_sold, sale_amt, status,
game_name, release_date, developer, publisher, platform_name, type
)
SELECT
gsf.sale_id,
gs.game_id,
gs.platform_id,
gs.sale_date,
gs.units_sold,
gs.sale_amt,
gs.status,
gd.game_name,
gd.release_date,
gd.developer,
gd.publisher,
pd.platform_name,
pd.type
FROM game_sales_fact AS gsf
JOIN game_sales FOR SYSTEM_TIME AS OF PROCTIME() AS gs
ON gsf.sale_id = gs.sale_id
JOIN game_dimension FOR SYSTEM_TIME AS OF PROCTIME() AS gd
ON gs.game_id = gd.game_id
JOIN platform_dimension FOR SYSTEM_TIME AS OF PROCTIME() AS pd
ON gs.platform_id = pd.platform_id;
ステップ 3:ジョブの開始
-
開発コンソールで、[O&M] > [デプロイメント] を選択し、両方のジョブデプロイメントを開始します。
-
両方のジョブが 実行中 状態になったら、HoloWeb に移動し、
game_sales_detailsテーブルをクエリします。SELECT * FROM game_sales_details;ステップ 1 でシードされた最初の行が結果に表示されます。
クエリは、以下のフィールド値を持つ 1 つのレコードを返します。
-
sale_id: 0 -
game_id: 0 -
platform_id: 101 -
sale_date: 2024-01-01 -
units_sold: 500 -
sale_amt: 2500 -
status: 1 -
game_name: SpaceInvaders -
release_date: 2023-06-15
-
ステップ 4:データの更新とクエリ
MongoDB の game_sales およびディメンションテーブルへの変更は、Hologres に自動的に伝播します。以下の例では、各更新タイプを実演します。
ファクトテーブルの更新
-
game_salesにさらに 5 行を挿入します。db.game_sales.insert( [ {sale_id:1,game_id:101,platform_id:1,"sale_date":"2024-01-01",units_sold:500,sale_amt:2500,status:1}, {sale_id:2,game_id:102,platform_id:2,"sale_date":"2024-08-02",units_sold:400,sale_amt:2000,status:1}, {sale_id:3,game_id:103,platform_id:1,"sale_date":"2024-08-03",units_sold:300,sale_amt:1500,status:1}, {sale_id:4,game_id:101,platform_id:3,"sale_date":"2024-08-04",units_sold:200,sale_amt:1000,status:1}, {sale_id:5,game_id:104,platform_id:2,"sale_date":"2024-08-05",units_sold:100,sale_amt:3000,status:1} ] );Hologres で
game_sales_detailsをクエリします。5 つの新しい行が表示されます。クエリ結果のテーブルには、MongoDB から同期されたフィールドに加えて、
game_name(SpaceInvaders、PuzzleQuest、RacingFever、AdventureLand など) とrelease_dateの列が含まれており、関連するゲームの詳細が表示されます。 -
sale_dateを2024-01-01から2024-08-01に更新します。db.game_sales.updateMany({"sale_date": "2024-01-01"}, {$set: {"sale_date": "2024-08-01"}});game_sales_detailsをクエリします。sale_date列に新しい値が反映されます。 -
statusを0に設定して、sale_id = 5の行を論理削除します。db.game_sales.updateMany({"sale_id": 5}, {$set: {"status": 0}});game_sales_detailsをクエリします。sale_id = 5のstatus列が0に変わります。
ディメンションテーブルの更新
-
ディメンションテーブルに新しいゲームとプラットフォームを追加します。
// 新しいゲーム db.game_dimension.insert( [ {game_id:105,"game_name":"HSHWK","release_date":"2024-08-20","developer":"GameSC","publisher":"GameSC"}, {game_id:106,"game_name":"HPBUBG","release_date":"2018-01-01","developer":"BLUE","publisher":"KK"} ] ); // 新しいプラットフォーム db.platform_dimension.insert( [ {platform_id:4,"platform_name":"Steam","type":"PC"}, {platform_id:5,"platform_name":"Epic","type":"PC"} ] );ディメンションテーブルへの挿入だけでは同期はトリガーされません。パイプラインは
game_salesの変更によって駆動されます。対応する販売レコードを挿入して、ワイドテーブルの更新をトリガーします。db.game_sales.insert( [ {sale_id:6,game_id:105,platform_id:4,"sale_date":"2024-09-01",units_sold:400,sale_amt:2000,status:1}, {sale_id:7,game_id:106,platform_id:1,"sale_date":"2024-09-01",units_sold:300,sale_amt:1500,status:1} ] );game_sales_detailsをクエリします。エンリッチされたディメンションデータを持つ 2 つの新しい行が表示されます。 -
MongoDB のディメンションデータを更新します。
// リリース日の更新 db.game_dimension.updateMany({"release_date": "2018-01-01"}, {$set: {"release_date": "2024-01-01"}}); // プラットフォームタイプの更新 db.platform_dimension.updateMany({"type": "PC"}, {$set: {"type": "Swich"}});更新されたフィールドは、Hologres の対応する行に伝播します。