All Products
Search
Document Center

Realtime Compute for Apache Flink:Analisis perilaku pengguna dengan Flink, MongoDB, dan Hologres

Last Updated:Aug 13, 2026

Pemrosesan data perilaku pengguna menantang karena volumenya yang besar dan formatnya yang beragam. Model tabel lebar tradisional menawarkan kinerja kueri yang efisien, tetapi dengan biaya redundansi data tinggi, overhead penyimpanan yang meningkat, serta siklus pemeliharaan yang lambat. Tutorial ini menjelaskan cara menggunakan Realtime Compute for Apache Flink, ApsaraDB for MongoDB, dan Hologres untuk membangun pipeline tabel lebar real-time yang mengatasi trade-off tersebut.

Cara kerja

Realtime Compute for Apache Flink menangani pemrosesan aliran. ApsaraDB for MongoDB menyimpan tabel fakta dan dimensi sebagai database NoSQL berbasis dokumen dengan skema fleksibel dan throughput baca/tulis yang tinggi. Hologres berfungsi sebagai gudang data analitik—data dapat dikueri segera setelah ditulis.

Pipeline ini menggunakan dua pekerjaan Flink yang dihubungkan oleh topik Kafka:

  1. Job 1 membaca aliran change data capture (CDC) dari MongoDB. Saat tabel fakta (game_sales) berubah, kunci primernya (PK) dialirkan langsung ke Kafka. Saat tabel dimensi berubah, lookup join mengidentifikasi catatan tabel fakta yang terpengaruh dan mengirim PK-nya ke Kafka.

  2. Job 2 mengonsumsi PK dari Kafka, melakukan lookup join terhadap tabel fakta dan dimensi MongoDB untuk merekonstruksi catatan tabel lebar lengkap, lalu melakukan upsert hasilnya ke Hologres.

image

Manfaat:

  • Throughput tulis tinggi: ApsaraDB for MongoDB menangani operasi baca/tulis konkurensi tinggi dalam kluster sharded, sehingga kinerja dan kapasitas penyimpanan dapat diskalakan untuk menampung volume besar pembaruan yang sering.

  • Propagasi perubahan efisien: Hanya PK dari catatan yang terpengaruh yang diteruskan ke Kafka, bukan seluruh baris. Hal ini meminimalkan pemrosesan ulang terlepas dari total volume data.

  • Kueri real-time: Hologres mendukung upsert dengan latensi rendah dan membuat data langsung dapat dikueri setelah setiap penulisan.

Praktik langsung

Pada akhir tutorial ini, Anda akan memiliki pipeline aktif di mana perubahan pada MongoDB—baik pada tabel fakta maupun tabel dimensi—secara otomatis dipropagasikan ke tabel lebar Hologres dan langsung dapat dikueri.

image

Pipeline ini menggabungkan tiga koleksi MongoDB menjadi satu tabel lebar 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>

Alur propagasi perubahan

Setiap perubahan mengalir melalui empat tahap:

  1. Capture: Mendeteksi perubahan real-time dari tabel dimensi MongoDB.

  2. Propagate: Saat tabel dimensi berubah, Job 1 melakukan lookup join (misalnya, pada game_id) untuk menemukan baris yang terpengaruh di tabel fakta dan mengekstraksi PK-nya (misalnya, sale_id).

  3. Trigger: Mengirim PK ke Kafka untuk memberi tahu Job 2 tentang refresh yang tertunda.

  4. Upsert: Job 2 mengambil data terbaru, merekonstruksi baris tabel lebar, lalu melakukan upsert ke Hologres.

Prasyarat

Sebelum memulai, pastikan Anda telah memiliki:

Langkah 1: Siapkan data

Buat koleksi MongoDB

  1. Masuk ke instans ApsaraDB for MongoDB Anda.

  2. Tambahkan Blok CIDR ruang kerja Flink Anda ke daftar putih MongoDB. Untuk detailnya, lihat Konfigurasikan daftar putih untuk instans dan Bagaimana cara mengonfigurasi daftar putih?

  3. Di editor SQL Konsol Data Management (DMS), buat database mongo_test:

    use mongo_test;
  4. Buat koleksi game_sales, game_dimension, dan platform_dimension, lalu masukkan data sampel:

    // Tabel penjualan game (status: 1 = aktif, 0 = dihapus secara logis)
    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},
      ]
    );
    
    // Tabel dimensi game
    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"},
      ]
    );
    
    // Tabel dimensi platform
    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. Verifikasi penyisipan:

    db.game_sales.find();
    db.game_dimension.find();
    db.platform_dimension.find();

    Kueri db.game_dimension.find() mengembalikan 4 catatan:

    • 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

Buat tabel Hologres

  1. Masuk ke Konsol Hologres, klik Instances di panel navigasi kiri, lalu klik instans Hologres Anda. Di pojok kanan atas, klik Connect to Instance.

  2. Di bilah navigasi atas, klik Metadata Management > Create Database. Masukkan test di bidang Database Name, atur Policy ke SPM, lalu klik OK. Untuk informasi lebih lanjut, lihat Buat database.

    Di kotak dialog, pilih instans bernama User-behavior-test, atur Log On Immediately ke Yes, lalu klik OK.

  3. Di bilah navigasi atas, klik SQL Editor. Klik ikon SQL untuk membuat kueri SQL baru, pilih instans dan database target, lalu jalankan pernyataan berikut untuk membuat tabel lebar 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)
    );

Buat topik Kafka

  1. Masuk ke Konsol ApsaraMQ for Kafka. Klik Instances di panel navigasi kiri, lalu klik instans Anda.

  2. Di panel navigasi kiri, klik Whitelist Management dan tambahkan Blok CIDR ruang kerja Flink Anda.

  3. Di panel navigasi kiri, klik Topics > Create Topic. Di panel kanan, masukkan game_sales_fact di bidang Name, masukkan deskripsi, dan pertahankan nilai default untuk semua bidang lainnya. Klik OK.

Langkah 2: Buat pekerjaan stream

Job 1: Tulis kunci primer ke Kafka

Job 1 memantau ketiga koleksi MongoDB. Saat game_sales berubah, sale_id dialirkan langsung ke Kafka. Saat tabel dimensi berubah, lookup join terhadap game_sales mengambil nilai sale_id yang terpengaruh, lalu dialirkan ke Kafka.

Konektor MongoDB menjalankan dua peran dalam pipeline ini. Di Job 1, ia berfungsi sebagai sumber CDC, membaca aliran perubahan MongoDB. Di Job 2, ia berfungsi sebagai sumber lookup, mengambil status terkini setiap dokumen berdasarkan kunci primer. Kedua peran menggunakan konfigurasi konektor yang sama.
image
  1. Masuk ke Konsol Realtime Compute for Apache Flink.

  2. Di kolom Actions ruang kerja Anda, klik Console.

  3. Di menu navigasi kiri, klik Development > ETL.

  4. Klik New Blank Stream Draft.

  5. Di dialog New Draft, masukkan dwd_mongo_kafka di Name, pilih versi engine, lalu klik Create.

  6. Salin SQL berikut ke editor. Masing-masing dari tiga pernyataan INSERT secara independen menangkap perubahan dari satu koleksi MongoDB dan mengalirkan nilai sale_id yang terpengaruh ke sink Kafka. Ini memastikan Hologres diperbarui secara akurat dan real-time setiap kali tabel berubah. Simpan nilai sensitif seperti string koneksi dan password sebagai variabel, bukan hardcoding. Untuk informasi lebih lanjut, lihat Kelola variabel.

    Waktu game_dimension state untuk game_id = 101
    T1 game_name = "SpaceInvaders"
    T2 Diperbarui ke game_name = "SpaceInvaders_v2"
    -- Sumber: game_sales (aliran CDC)
    CREATE TEMPORARY TABLE game_sales
    (
      `_id`       STRING,    -- ID yang dihasilkan otomatis oleh MongoDB
      sale_id     INT,       -- ID Penjualan
      PRIMARY KEY (_id) NOT ENFORCED
    )
    WITH (
      'connector' = 'mongodb',
      'uri' = '${secret_values.MongoDB-URI}',
      'database' = 'mongo_test',
      'collection' = 'game_sales'
    );
    
    -- Sumber: game_dimension (aliran 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'
    );
    
    -- Sumber: platform_dimension (aliran 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'
    );
    
    -- Sumber lookup: game_sales (digunakan untuk join perubahan dimensi)
    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: Topik Kafka yang menyimpan PK yang terpengaruh
    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'  -- Diperlukan saat menulis ke ApsaraMQ for Kafka
    );
    
    BEGIN STATEMENT SET;
    
    -- Alirkan PK dari perubahan game_sales
    INSERT INTO game_sales_fact (sale_id)
    SELECT sale_id FROM game_sales;
    
    -- Alirkan PK dari baris game_sales yang terpengaruh oleh perubahan game_dimension
    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;
    
    -- Alirkan PK dari baris game_sales yang terpengaruh oleh perubahan platform_dimension
    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;

    Tentang lookup join Klausul FOR SYSTEM_TIME AS OF PROCTIME() mendefinisikan lookup join (jenis join temporal). Pada saat baris sumber diproses, join mengambil baris tabel dimensi yang sesuai dari MongoDB dan membekukan snapshot tersebut untuk hasil join. Jika tabel dimensi diperbarui kemudian, baris yang sudah diproses tidak terpengaruh. Contohnya: Baris game_sales yang diproses pada T1 melakukan join dengan "SpaceInvaders". Baris yang diproses pada T2 melakukan join dengan "SpaceInvaders_v2". Kondisi join adalah gd.game_id = gs.game_id dan pd.platform_id = gs.platform_id. Untuk informasi lebih lanjut, lihat Pernyataan JOIN untuk tabel dimensi dan Pilih Kafka, Upsert Kafka, atau katalog Kafka JSON.

  7. Di pojok kanan atas, klik Deploy. Di dialog, klik Confirm. Untuk informasi lebih lanjut, lihat Terapkan pekerjaan.

Job 2: Rekonstruksi tabel lebar dan upsert ke Hologres

Job 2 mengonsumsi nilai sale_id dari topik Kafka game_sales_fact, melakukan lookup join terhadap tabel fakta dan dimensi MongoDB, lalu melakukan upsert baris tabel lebar hasilnya ke Hologres.

image

Ikuti langkah-langkah di Job 1 untuk membuat draft baru bernama dws_kafka_mongo_holo dan terapkan dengan SQL berikut:

-- Sumber: Topik Kafka yang menyediakan PK yang terpengaruh
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'
);

-- Sumber lookup: Tabel fakta 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'
);

-- Sumber lookup: 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'
);

-- Sumber lookup: 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: Tabel lebar 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',       -- Hanya insert atau update; jangan pernah menghapus baris
  'sink.on-conflict-action' = 'INSERT_OR_UPDATE',  -- Aktifkan pembaruan kolom parsial
  '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;

Langkah 3: Jalankan pekerjaan

  1. Di Konsol Development, pilih O&M > Deployments dan jalankan kedua deployment pekerjaan.

  2. Setelah kedua pekerjaan mencapai status Running, buka HoloWeb dan kueri tabel game_sales_details:

    SELECT * FROM game_sales_details;

    Baris awal yang disisipkan di Langkah 1 muncul dalam hasil.

    Kueri mengembalikan satu catatan dengan nilai bidang berikut:

    • 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

Langkah 4: Perbarui dan kueri data

Perubahan pada game_sales dan tabel dimensi di MongoDB dipropagasikan ke Hologres secara otomatis. Contoh berikut menunjukkan setiap jenis pembaruan.

Pembaruan tabel fakta

  1. Sisipkan lima baris lagi ke game_sales:

    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}
      ]
    );

    Kueri game_sales_details di Hologres. Lima baris baru muncul.

    Tabel hasil kueri berisi kolom game_name (seperti SpaceInvaders, PuzzleQuest, RacingFever, dan AdventureLand) dan kolom release_date selain bidang yang disinkronkan dari MongoDB, menampilkan detail game terkait.

  2. Perbarui sale_date dari 2024-01-01 menjadi 2024-08-01:

    db.game_sales.updateMany({"sale_date": "2024-01-01"}, {$set: {"sale_date": "2024-08-01"}});

    Kueri game_sales_details. Kolom sale_date mencerminkan nilai baru.

  3. Hapus secara logis baris dengan sale_id = 5 dengan mengatur status ke 0:

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

    Kueri game_sales_details. Kolom status untuk sale_id = 5 berubah menjadi 0.

Pembaruan tabel dimensi

  1. Tambahkan game dan platform baru ke tabel dimensi:

    // Game baru
    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"}
      ]
    );
    
    // Platform baru
    db.platform_dimension.insert(
      [
    {platform_id:4,"platform_name":"Steam","type":"PC"},
    {platform_id:5,"platform_name":"Epic","type":"PC"}
      ]
    );

    Menyisipkan ke tabel dimensi saja tidak memicu sinkronisasi—pipeline didorong oleh perubahan pada game_sales. Sisipkan catatan penjualan terkait untuk memicu pembaruan tabel lebar:

    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}
      ]
    );

    Kueri game_sales_details. Dua baris baru muncul dengan data dimensi yang diperkaya.

  2. Perbarui data dimensi di MongoDB:

    // Perbarui tanggal rilis
    db.game_dimension.updateMany({"release_date": "2018-01-01"}, {$set: {"release_date": "2024-01-01"}});
    
    // Perbarui tipe platform
    db.platform_dimension.updateMany({"type": "PC"}, {$set: {"type": "Swich"}});

    Bidang yang diperbarui dipropagasikan ke baris terkait di Hologres.

Lanjutan