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:
-
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. -
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.
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.
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:
-
Capture: Mendeteksi perubahan real-time dari tabel dimensi MongoDB.
-
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). -
Trigger: Mengirim PK ke Kafka untuk memberi tahu Job 2 tentang refresh yang tertunda.
-
Upsert: Job 2 mengambil data terbaru, merekonstruksi baris tabel lebar, lalu melakukan upsert ke Hologres.
Prasyarat
Sebelum memulai, pastikan Anda telah memiliki:
-
Ruang kerja Realtime Compute for Apache Flink yang menjalankan VVR 8.0.5 atau lebih baru. Untuk informasi lebih lanjut, lihat Aktifkan Realtime Compute for Apache Flink.
-
Instans ApsaraDB for MongoDB yang menjalankan versi 4.0 atau lebih baru. Untuk informasi lebih lanjut, lihat Buat instans kluster sharded.
-
Instans eksklusif Hologres yang menjalankan versi 1.3 atau lebih baru. Untuk informasi lebih lanjut, lihat Beli instans Hologres.
-
Instans ApsaraMQ for Kafka. Untuk informasi lebih lanjut, lihat Terapkan instans ApsaraMQ for Kafka.
-
Keempat instans berada dalam Virtual Private Cloud (VPC) yang sama. Jika berada di VPC berbeda, bangun konektivitas cross-VPC atau aktifkan akses Internet untuk Realtime Compute for Apache Flink. Untuk informasi lebih lanjut, lihat Bagaimana Realtime Compute for Apache Flink mengakses layanan lintas VPC? dan Bagaimana Realtime Compute for Apache Flink mengakses Internet?
-
Izin Pengguna RAM atau Peran RAM untuk sumber daya terkait.
Langkah 1: Siapkan data
Buat koleksi MongoDB
-
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?
-
Di editor SQL Konsol Data Management (DMS), buat database
mongo_test:use mongo_test; -
Buat koleksi
game_sales,game_dimension, danplatform_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"} ] ); -
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
-
Masuk ke Konsol Hologres, klik Instances di panel navigasi kiri, lalu klik instans Hologres Anda. Di pojok kanan atas, klik Connect to Instance.
-
Di bilah navigasi atas, klik Metadata Management > Create Database. Masukkan
testdi 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.
-
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
-
Masuk ke Konsol ApsaraMQ for Kafka. Klik Instances di panel navigasi kiri, lalu klik instans Anda.
-
Di panel navigasi kiri, klik Whitelist Management dan tambahkan Blok CIDR ruang kerja Flink Anda.
-
Di panel navigasi kiri, klik Topics > Create Topic. Di panel kanan, masukkan
game_sales_factdi 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.
-
Masuk ke Konsol Realtime Compute for Apache Flink.
-
Di kolom Actions ruang kerja Anda, klik Console.
-
Di menu navigasi kiri, klik Development > ETL.
-
Klik New Blank Stream Draft.
-
Di dialog New Draft, masukkan
dwd_mongo_kafkadi Name, pilih versi engine, lalu klik Create. -
Salin SQL berikut ke editor. Masing-masing dari tiga pernyataan
INSERTsecara independen menangkap perubahan dari satu koleksi MongoDB dan mengalirkan nilaisale_idyang 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_dimensionstate untukgame_id = 101T1 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: Barisgame_salesyang diproses pada T1 melakukan join dengan"SpaceInvaders". Baris yang diproses pada T2 melakukan join dengan"SpaceInvaders_v2". Kondisi join adalahgd.game_id = gs.game_iddanpd.platform_id = gs.platform_id. Untuk informasi lebih lanjut, lihat Pernyataan JOIN untuk tabel dimensi dan Pilih Kafka, Upsert Kafka, atau katalog Kafka JSON. -
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.
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
-
Di Konsol Development, pilih O&M > Deployments dan jalankan kedua deployment pekerjaan.
-
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
-
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_detailsdi Hologres. Lima baris baru muncul.Tabel hasil kueri berisi kolom
game_name(seperti SpaceInvaders, PuzzleQuest, RacingFever, dan AdventureLand) dan kolomrelease_dateselain bidang yang disinkronkan dari MongoDB, menampilkan detail game terkait. -
Perbarui
sale_datedari2024-01-01menjadi2024-08-01:db.game_sales.updateMany({"sale_date": "2024-01-01"}, {$set: {"sale_date": "2024-08-01"}});Kueri
game_sales_details. Kolomsale_datemencerminkan nilai baru. -
Hapus secara logis baris dengan
sale_id = 5dengan mengaturstatuske0:db.game_sales.updateMany({"sale_id": 5}, {$set: {"status": 0}});Kueri
game_sales_details. Kolomstatusuntuksale_id = 5berubah menjadi0.
Pembaruan tabel dimensi
-
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. -
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.