すべてのプロダクト
Search
ドキュメントセンター

Realtime Compute for Apache Flink:Flink、MongoDB、Hologres を使用したユーザー行動分析

最終更新日:Aug 13, 2026

ユーザー行動データの処理は、その膨大な量と多様なフォーマットのため、困難を伴います。従来のワイドテーブルモデルは、効率的なクエリパフォーマンスを提供する一方で、高いデータ冗長性、ストレージオーバーヘッドの増加、メンテナンスサイクルの長期化といった代償を伴います。このチュートリアルでは、Realtime Compute for Apache Flink、ApsaraDB for MongoDB、Hologres を使用して、これらのトレードオフに対応するリアルタイムのワイドテーブルパイプラインを構築する方法を説明します。

仕組み

Realtime Compute for Apache Flink はストリーム処理を担います。ApsaraDB for MongoDB は、柔軟なスキーマと高い読み書きスループットを持つドキュメント指向の NoSQL データベースとして、ファクトテーブルとディメンションテーブルを格納します。Hologres は分析用データウェアハウスとして機能し、データは書き込み後すぐにクエリ可能になります。

このパイプラインは、Kafka トピックで接続された 2 つの Flink ジョブを使用します。

  1. ジョブ 1 は MongoDB からの変更データキャプチャ (CDC) ストリームを読み取ります。ファクトテーブル (game_sales) が変更されると、そのプライマリキー (PK) は直接 Kafka にストリーミングされます。ディメンションテーブルが変更されると、ルックアップ結合によって影響を受けるファクトテーブルのレコードが特定され、その PK が Kafka に送信されます。

  2. ジョブ 2 は Kafka から PK を消費し、MongoDB のファクトテーブルとディメンションテーブルに対してルックアップ結合を実行して完全なワイドテーブルレコードを再構築し、その結果を Hologres にアップサートします。

image

メリット:

  • 高い書き込みスループット: ApsaraDB for MongoDB は、シャードクラスターでの高同時実行数の読み書きを処理し、パフォーマンスとストレージをスケーリングして、大量の頻繁な更新に対応します。

  • 効率的な変更の伝播: 影響を受けるレコードの PK のみが Kafka に転送され、行全体は転送されません。これにより、総データ量に関係なく、再処理が最小限に抑えられます。

  • リアルタイムクエリ: Hologres は低レイテンシーのアップサートをサポートし、各書き込みの直後にデータをクエリ可能にします。

ハンズオン

このチュートリアルを完了すると、ファクトテーブルとディメンションテーブルの両方に対する MongoDB の変更が、Hologres のワイドテーブルに自動的に伝播し、即座にクエリ可能になる稼働中のパイプラインが完成します。

image

このパイプラインは、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 つのステージを経て流れます。

  1. キャプチャ: MongoDB ディメンションテーブルからのリアルタイムの変更を検出します。

  2. 伝播: ディメンションテーブルが変更されると、ジョブ 1 はルックアップ結合 (例えば、game_id で) を実行してファクトテーブル内の影響を受ける行を見つけ、その PK (例えば、sale_id) を抽出します。

  3. トリガー: PK を Kafka に送信して、保留中のリフレッシュをジョブ 2 に通知します。

  4. アップサート: ジョブ 2 は最新のデータをフェッチし、ワイドテーブルの行を再構築して、Hologres にアップサートします。

前提条件

開始する前に、以下が揃っていることを確認してください。

ステップ 1:データの準備

MongoDB コレクションの作成

  1. ご利用の ApsaraDB for MongoDB インスタンスにログインします。

  2. ご利用の Flink ワークスペースの CIDR ブロックを MongoDB のホワイトリストに追加します。詳細については、「インスタンスのホワイトリストの設定」および「ホワイトリストの設定方法」をご参照ください。

  3. Data Management (DMS) コンソールの SQL エディターで、mongo_test データベースを作成します。

    use mongo_test;
  4. 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"}
      ]
    );
  5. 挿入を検証します。

    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 テーブルの作成

  1. Hologres コンソールにログインし、左側のナビゲーションウィンドウで [インスタンス] をクリックし、ご利用の Hologres インスタンスをクリックします。右上隅の [インスタンスに接続] をクリックします。

  2. 上部のナビゲーションバーで、[メタデータ管理] > [データベースの作成] をクリックします。[データベース名] フィールドに test と入力し、[ポリシー] を [SPM] に設定し、[OK] をクリックします。詳細については、「データベースの作成」をご参照ください。

    ダイアログボックスで、インスタンス [User-behavior-test] を選択し、[すぐにログオン] を [はい] に設定し、[OK] をクリックします。

  3. 上部のナビゲーションバーで、[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 トピックの作成

  1. ApsaraMQ for Kafka コンソールにログインします。左側のナビゲーションウィンドウで [インスタンス] をクリックし、ご利用のインスタンスをクリックします。

  2. 左側のナビゲーションウィンドウで、[ホワイトリスト管理] をクリックし、ご利用の Flink ワークスペースの CIDR ブロックを追加します。

  3. 左側のナビゲーションウィンドウで、[トピック] > [トピックの作成] をクリックします。右側のペインで、[名前] フィールドに 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 では、プライマリキーによって各ドキュメントの現在の状態をフェッチするルックアップソースとして機能します。両方の役割で同じコネクタ構成が使用されます。
image
  1. Realtime Compute for Apache Flink コンソールにログインします。

  2. ご利用のワークスペースの [操作] 列で、[コンソール] をクリックします。

  3. 左側のナビゲーションメニューで、[開発] > [ETL] をクリックします。

  4. [新しい空白のストリームドラフト] をクリックします。

  5. [新しいドラフト] ダイアログで、[名前] に dwd_mongo_kafka と入力し、エンジンバージョンを選択して、[作成] をクリックします。

  6. 次の SQL をエディターにコピーします。3 つの INSERT 文はそれぞれ、独立して 1 つの MongoDB コレクションからの変更をキャプチャし、影響を受ける sale_id の値を Kafka sink にストリーミングします。これにより、いずれかのテーブルが変更されたときに、Hologres が正確かつリアルタイムに更新されることが保証されます。接続文字列やパスワードなどの機密性の高い値は、ハードコーディングするのではなく、変数として格納してください。詳細については、「変数の管理」をご参照ください。

    時間 game_id = 101 の game_dimension
    T1 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 カタログの選択」をご参照ください。

  7. 右上隅で [デプロイ] をクリックします。ダイアログで [確認] をクリックします。詳細については、「ジョブのデプロイ」をご参照ください。

ジョブ 2:ワイドテーブルの再構築と Hologres へのアップサート

ジョブ 2 は game_sales_fact Kafka トピックから sale_id の値を消費し、MongoDB のファクトテーブルとディメンションテーブルに対してルックアップ結合を実行し、結果のワイドテーブル行を Hologres にアップサートします。

image

ジョブ 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:ジョブの開始

  1. 開発コンソールで、[O&M] > [デプロイメント] を選択し、両方のジョブデプロイメントを開始します。

  2. 両方のジョブが 実行中 状態になったら、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 に自動的に伝播します。以下の例では、各更新タイプを実演します。

ファクトテーブルの更新

  1. 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 の列が含まれており、関連するゲームの詳細が表示されます。

  2. 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 列に新しい値が反映されます。

  3. status を 0 に設定して、sale_id = 5 の行を論理削除します。

    db.game_sales.updateMany({"sale_id": 5}, {$set: {"status": 0}});

    game_sales_details をクエリします。sale_id = 5 の status 列が 0 に変わります。

ディメンションテーブルの更新

  1. ディメンションテーブルに新しいゲームとプラットフォームを追加します。

    // 新しいゲーム
    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 つの新しい行が表示されます。

  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 の対応する行に伝播します。

次のステップ