Topik ini menjelaskan cara membangun danau data terpadu streaming menggunakan Data Lake Formation (DLF), Realtime Compute for Apache Flink, dan StarRocks.
Informasi latar belakang
Seiring percepatan transformasi digital bisnis, permintaan terhadap data yang tepat waktu semakin meningkat. Gudang data offline tradisional mengandalkan pekerjaan penjadwalan berbasis waktu untuk menambahkan data secara inkremental ke lapisan Operational Data Store (ODS), Data Warehouse Detail (DWD), Data Warehouse Service (DWS), dan Application Data Service (ADS). Model ini menghadapi dua tantangan utama. Pertama, latensi tinggi. Siklus penjadwalan biasanya berada pada skala jam atau hari, sehingga konsumen data tidak dapat memperoleh informasi secara real-time. Kedua, biaya tinggi. Pembaruan data bergantung pada penggantian partisi, yang mengharuskan sistem membaca ulang data asli untuk menggabungkan perubahan—mengonsumsi banyak sumber daya komputasi.
Membangun danau data terpadu streaming berbasis Flink dan DLF merupakan solusi efektif untuk mengatasi tantangan tersebut. Flink memungkinkan aliran data real-time antar lapisan gudang data, sedangkan DLF menggunakan mekanisme pembaruan data yang efisien sehingga konsumen downstream dapat memperoleh perubahan data terbaru dengan latensi tingkat menit. Arsitektur ini secara signifikan mengurangi latensi data serta biaya penyimpanan dan komputasi.
Untuk informasi lebih lanjut tentang fitur produk, lihat Manfaat.
Arsitektur dan manfaat
Arsitektur
Flink adalah mesin komputasi aliran andal yang mendukung pemrosesan data real-time dalam jumlah besar. DLF menggunakan Paimon sebagai format penyimpanan danau terpadu untuk pemrosesan aliran dan batch. Paimon mendukung pembaruan throughput tinggi dan kueri latensi rendah, serta terintegrasi secara mendalam dengan Flink untuk menyediakan solusi danau data terpadu streaming yang terintegrasi.
Flink menulis data dari sumber data ke Paimon untuk membentuk lapisan ODS.
Flink berlangganan changelog data dari lapisan ODS, memprosesnya, lalu menulis kembali ke Paimon untuk membentuk lapisan DWD.
Flink berlangganan changelog data dari lapisan DWD, memprosesnya, lalu menulis kembali ke Paimon untuk membentuk lapisan DWS.
Terakhir, StarRocks di E-MapReduce open source membaca tabel eksternal Paimon untuk menyediakan layanan kueri bagi aplikasi.

Manfaat
Solusi ini memberikan manfaat berikut:
Setiap lapisan data Paimon dapat menyebarkan perubahan ke sistem downstream dengan latensi tingkat menit, mengurangi latensi gudang data offline tradisional dari hitungan jam atau hari menjadi hitungan menit.
Setiap lapisan data Paimon dapat langsung menerima data perubahan tanpa perlu mengganti partisi, sehingga secara signifikan mengurangi biaya pembaruan dan koreksi data pada gudang data offline tradisional serta mengatasi kesulitan dalam melakukan kueri, pembaruan, dan koreksi data pada lapisan antara.
Modelnya terpadu dan arsitekturnya disederhanakan. Logika pipeline ekstrak, transformasi, dan muat (ETL) diimplementasikan berdasarkan Flink SQL, sedangkan data pada lapisan ODS, DWD, dan DWS disimpan secara seragam dalam Paimon. Hal ini mengurangi kompleksitas arsitektur dan meningkatkan efisiensi pemrosesan data.
Solusi ini mengandalkan tiga kemampuan inti Paimon. Tabel berikut menjelaskan detailnya.
Kemampuan inti Paimon | Rincian |
Pembaruan tabel primary key | Paimon menggunakan struktur data Log-structured Merge-tree (LSM) pada lapisan dasar untuk mencapai pembaruan data yang efisien. Untuk informasi lebih lanjut tentang tabel primary key Paimon dan struktur data dasar Paimon, lihat Primary Key Table dan File Layouts. |
Changelog producer | Paimon dapat menghasilkan data changelog lengkap untuk setiap aliran data input. Semua data update_after memiliki data update_before yang sesuai. Hal ini memastikan bahwa perubahan data dapat sepenuhnya diteruskan ke sistem downstream. Untuk informasi lebih lanjut, lihat Changelog producer. |
Mesin merge | Ketika tabel primary key Paimon menerima beberapa data dengan primary key yang sama, tabel sink Paimon menggabungkan data tersebut menjadi satu untuk menjaga keunikan primary key. Paimon mendukung berbagai perilaku penggabungan data, seperti deduplikasi, pembaruan parsial, dan pra-agregasi. Untuk informasi lebih lanjut, lihat Merge engine. |
Praktik terbaik
Contoh ini menunjukkan cara membangun danau data terpadu streaming untuk platform e-commerce guna memproses dan membersihkan data serta melakukan kueri data dari aplikasi lapisan atas. Dengan demikian, data dilapisankan dan digunakan kembali untuk mendukung berbagai skenario bisnis, seperti kueri laporan analisis data dashboard transaksi, analisis data perilaku, pelabelan profil pengguna, dan rekomendasi personalisasi.

Bangun lapisan ODS: ingest data secara real-time dari database bisnis ke gudang data.Database MySQL berisi tabel bisnis berikut: orders, orders_pay, dan product_catalog. Realtime Compute for Apache Flink menulis data dari tabel-tabel ini ke Object Storage Service (OSS) secara real-time dan menyimpannya dalam format Apache Paimon untuk membentuk lapisan ODS.
Bangun lapisan DWD: tabel lebar.Realtime Compute for Apache Flink menggunakan mekanisme penggabungan data pembaruan parsial Apache Paimon untuk menggabungkan data dari tabel orders, orders_pay, dan product_catalog menjadi tabel lebar pada lapisan DWD dan menghasilkan changelog dengan latensi tingkat menit.
Bangun lapisan DWS: perhitungan metrik.Realtime Compute for Apache Flink mengonsumsi changelog dari tabel lebar secara real-time dan menggunakan mekanisme penggabungan data agregasi Apache Paimon untuk menghasilkan tabel antara agregat pengguna-pedagang bernama dwm_users_shops pada lapisan tengah gudang data (DWM), lalu akhirnya menghasilkan tabel metrik agregat pengguna bernama dws_users dan tabel metrik agregat pedagang bernama dws_shops pada lapisan DWS.
Prasyarat
DLF telah diaktifkan. Untuk informasi lebih lanjut, lihat Memulai dengan DLF.
Realtime Compute for Apache Flink telah diaktifkan. Untuk informasi lebih lanjut, lihat Aktifkan Realtime Compute for Apache Flink.
EMR Serverless StarRocks telah diaktifkan. Untuk informasi lebih lanjut, lihat Memulai dengan instans Serverless StarRocks yang menggunakan arsitektur pemisahan penyimpanan-komputasi.
Instans StarRocks dan DLF harus berada di Wilayah yang sama dengan ruang kerja Flink.
Batasan
Hanya Ververica Runtime (VVR) 11.1.0 dan versi yang lebih baru yang mendukung solusi danau data terpadu streaming ini.
Bangun danau data terpadu streaming
Persiapkan sumber data CDC MySQL
Pada contoh ini, tiga tabel bisnis dibuat dalam database bernama order_dw pada instans ApsaraDB RDS for MySQL, dan data dimasukkan ke dalam tabel-tabel tersebut.
Buat instans ApsaraDB RDS for MySQL.
CatatanJika ApsaraDB RDS for MySQL dan ruang kerja Flink Anda tidak berada dalam VPC yang sama, lihat Bagaimana cara mengakses layanan lain lintas VPC?.
Buat database bernama
order_dwdan buat akun istimewa atau akun standar yang memiliki izin baca dan tulis pada databaseorder_dw.Buat tiga tabel dan masukkan data ke dalam tabel-tabel tersebut.
CREATE TABLE `orders` ( order_id bigint not null primary key, user_id varchar(50) not null, shop_id bigint not null, product_id bigint not null, buy_fee bigint not null, create_time timestamp not null, update_time timestamp not null default now(), state int not null ); CREATE TABLE `orders_pay` ( pay_id bigint not null primary key, order_id bigint not null, pay_platform int not null, create_time timestamp not null ); CREATE TABLE `product_catalog` ( product_id bigint not null primary key, catalog_name varchar(50) not null ); -- Persiapkan data. INSERT INTO product_catalog VALUES(1, 'phone_aaa'),(2, 'phone_bbb'),(3, 'phone_ccc'),(4, 'phone_ddd'),(5, 'phone_eee'); INSERT INTO orders VALUES (100001, 'user_001', 12345, 1, 5000, '2023-02-15 16:40:56', '2023-02-15 18:42:56', 1), (100002, 'user_002', 12346, 2, 4000, '2023-02-15 15:40:56', '2023-02-15 18:42:56', 1), (100003, 'user_003', 12347, 3, 3000, '2023-02-15 14:40:56', '2023-02-15 18:42:56', 1), (100004, 'user_001', 12347, 4, 2000, '2023-02-15 13:40:56', '2023-02-15 18:42:56', 1), (100005, 'user_002', 12348, 5, 1000, '2023-02-15 12:40:56', '2023-02-15 18:42:56', 1), (100006, 'user_001', 12348, 1, 1000, '2023-02-15 11:40:56', '2023-02-15 18:42:56', 1), (100007, 'user_003', 12347, 4, 2000, '2023-02-15 10:40:56', '2023-02-15 18:42:56', 1); INSERT INTO orders_pay VALUES (2001, 100001, 1, '2023-02-15 17:40:56'), (2002, 100002, 1, '2023-02-15 17:40:56'), (2003, 100003, 0, '2023-02-15 17:40:56'), (2004, 100004, 0, '2023-02-15 17:40:56'), (2005, 100005, 0, '2023-02-15 18:40:56'), (2006, 100006, 0, '2023-02-15 18:40:56'), (2007, 100007, 0, '2023-02-15 18:40:56');
Kelola metadata
Buat katalog Paimon
Login ke Konsol Realtime Compute for Apache Flink.
Pada bilah navigasi di sebelah kiri, pilih Metadata Management dan klik Create Catalog.
Pada tab Built-in Catalog, klik Apache Paimon, lalu klik Next.
Masukkan parameter berikut, pilih DLF sebagai kelas penyimpanan, lalu klik OK.
Item konfigurasi
Deskripsi
Wajib
Keterangan
metastore
Jenis metastore.
Ya
Contoh ini menggunakan kelas penyimpanan dlf.
catalog name
Nama katalog data DLF.
PentingJika Anda menggunakan Pengguna Resource Access Management (RAM) atau role, pastikan Anda memiliki izin baca dan tulis pada data DLF. Untuk informasi lebih lanjut, lihat Kelola izin data.
Ya
Pilih katalog data DLF yang sudah ada. Untuk membuat katalog data, lihat Katalog data.
Pada contoh ini, katalog data bernama
paimoncatalogtelah dibuat sebelumnya.Buat database order_dw yang sesuai dalam katalog data untuk menyinkronkan semua data tabel dari database order_dw di MySQL.
Pada bilah navigasi di sebelah kiri, pilih , lalu klik New untuk membuat kueri sementara.
-- Gunakan sumber data paimoncatalog USE CATALOG paimoncatalog; -- Buat database order_dw CREATE DATABASE order_dw;Pesan
The following statement has been executed successfully!menunjukkan bahwa database berhasil dibuat.
Untuk informasi lebih lanjut tentang cara menggunakan katalog Paimon, lihat Kelola katalog Paimon.
Buat katalog MySQL
Pada halaman Metadata Management, klik Create Catalog.
Pada tab Built-in Catalog, klik MySQL, lalu klik Next.
Masukkan parameter berikut dan klik OK untuk membuat katalog MySQL bernama mysqlcatalog.
Item konfigurasi
Deskripsi
Wajib
Keterangan
catalog name
Nama katalog.
Ya
Masukkan nama kustom dalam bahasa Inggris. Topik ini menggunakan mysqlcatalog sebagai contoh.
hostname
Alamat IP atau hostname database MySQL.
Ya
Untuk informasi lebih lanjut, lihat Lihat dan kelola titik akhir serta nomor port instans. Karena instans ApsaraDB RDS for MySQL dan Flink yang sepenuhnya dikelola berada dalam VPC yang sama, masukkan alamat jaringan internal.
port
Nomor port layanan database MySQL. Nilai default adalah 3306.
Tidak
Untuk informasi lebih lanjut, lihat Lihat dan kelola titik akhir serta nomor port instans.
default-database
Nama database MySQL default.
Ya
Topik ini menggunakan nama database yang akan disinkronkan, yaitu order_dw.
username
Username untuk layanan database MySQL.
Ya
Ini adalah akun yang dibuat di Persiapkan sumber data CDC MySQL.
password
Password untuk layanan database MySQL.
Ya
Ini adalah password untuk akun yang dibuat di Persiapkan sumber data CDC MySQL.
Bangun lapisan ODS: Ingesti data secara real-time ke gudang data
Gunakan Flink CDC untuk mengingesti data dari MySQL ke Paimon dan membangun lapisan ODS.
Buat pekerjaan ingest data.
Login ke Konsol Manajemen Realtime Compute for Apache Flink. Klik Console pada kolom Actions ruang kerja Anda untuk masuk ke Konsol Pengembangan.
Pada menu navigasi kiri, pilih . Klik . Pada dialog New Draft, masukkan
odspada bidang Name lalu klik Create. Pada editor SQL, salin dan tempel kode berikut:source: type: mysql name: MySQL Source hostname: rm-bp1e********566g.mysql.rds.aliyuncs.com port: 3306 username: ${secret_values.username} password: ${secret_values.password} tables: order_dw.\.* # Baca semua tabel di database MySQL order_dw. sink: type: paimon name: Paimon Sink catalog.properties.metastore: rest catalog.properties.uri: http://ap-southeast-5-vpc.dlf.aliyuncs.com catalog.properties.warehouse: paimoncatalog catalog.properties.token.provider: dlf pipeline: name: MySQL to Paimon PipelineItem konfigurasi
Deskripsi
Wajib?
Contoh
catalog.properties.metastoreJenis metastore. Atur ke
rest.Ya
rest
catalog.properties.token.providerPenyedia token. Atur ke
dlf.Ya
dlf
catalog.properties.uriURI server DLF. Format:
http://[region-id]-vpc.dlf.aliyuncs.com. Ganti[region-id]dengan ID wilayah aktual Anda (Lihat Endpoints).Ya
http://ap-southeast-5-vpc.dlf.aliyuncs.com
catalog.properties.warehouseNama katalog DLF.
Ya
paimoncatalog
Anda dapat mengonfigurasi properti tabel Paimon untuk meningkatkan performa penulisan. Untuk detailnya, lihat Optimalisasi performa.
Pada pojok kanan atas editor SQL, klik Deploy.
Pada menu navigasi kiri, pilih . Pada halaman Deployments, temukan deployment pekerjaan
odslalu klik Start pada kolom Actions. Pada panel Start Job, pilih Initial Mode lalu klik Start.
Lihat data yang disinkronkan dari MySQL ke Paimon.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Klik . Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu klik Run:
SELECT * FROM paimoncatalog.order_dw.orders ORDER BY order_id;
Bangun lapisan DWD: Tabel lebar
Buat tabel lebar bernama dwd_orders pada lapisan DWD di Apache Paimon.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada tab Scripts, klik . Pada editor SQL, salin kode berikut, pilih kodenya, lalu klik Run:
CREATE TABLE paimoncatalog.order_dw.dwd_orders ( order_id BIGINT, order_user_id STRING, order_shop_id BIGINT, order_product_id BIGINT, order_product_catalog_name STRING, order_fee BIGINT, order_create_time TIMESTAMP, order_update_time TIMESTAMP, order_state INT, pay_id BIGINT, pay_platform INT COMMENT 'platform 0: phone, 1: pc', pay_create_time TIMESTAMP, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'merge-engine' = 'partial-update', -- Gunakan mekanisme penggabungan data pembaruan parsial untuk menghasilkan tabel lebar. 'changelog-producer' = 'lookup' -- Gunakan mekanisme pembuatan data inkremental lookup untuk menghasilkan changelog dengan latensi rendah. );Jika pesan
Query has been executeddikembalikan, berarti tabel berhasil dibuat.Konsumsi changelog dari tabel orders dan orders_pay pada lapisan ODS secara real-time.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada halaman yang muncul, klik untuk membuat draft bernama
dwd. Pada editor SQL, salin dan tempel kode berikut, lalu deploy draft tersebut. Mulai deployment pekerjaan dengan status awal.Pekerjaan ini melakukan join tabel
ordersdengan tabel dimensi bernamaproduct_catalogdan menulis hasil join serta data dari tabelorders_payke tabel lebar bernamadwd_orders. Dalam proses ini, mekanisme penggabungan data pembaruan parsial Apache Paimon digunakan untuk menggabungkan data yang memiliki nilaiorder_idyang sama pada tabel orders danorders_pay.SET 'execution.checkpointing.max-concurrent-checkpoints' = '3'; SET 'table.exec.sink.upsert-materialize' = 'NONE'; SET 'execution.checkpointing.interval' = '10s'; SET 'execution.checkpointing.min-pause' = '10s'; -- Apache Paimon tidak mengizinkan Anda menggunakan beberapa pernyataan INSERT untuk menulis data ke tabel yang sama dalam deployment yang sama. Oleh karena itu, UNION ALL digunakan dalam contoh ini. INSERT INTO paimoncatalog.order_dw.dwd_orders SELECT o.order_id, o.user_id, o.shop_id, o.product_id, dim.catalog_name, o.buy_fee, o.create_time, o.update_time, o.state, NULL, NULL, NULL FROM paimoncatalog.order_dw.orders o LEFT JOIN paimoncatalog.order_dw.product_catalog FOR SYSTEM_TIME AS OF proctime() AS dim ON o.product_id = dim.product_id UNION ALL SELECT order_id, NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, pay_id, pay_platform, create_time FROM paimoncatalog.order_dw.orders_pay;Lihat data tabel lebar bernama dwd_orders.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada tab Scripts, klik . Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu Run:
SELECT * FROM paimoncatalog.order_dw.dwd_orders ORDER BY order_id;
Bangun lapisan DWS: Perhitungan metrik
Buat tabel metrik agregat bernama dws_users dan dws_shops pada lapisan DWS.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada tab Scripts, klik . Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu Run:
-- Buat tabel metrik agregat berdimensi pengguna. CREATE TABLE paimoncatalog.order_dw.dws_users ( user_id STRING, ds STRING, payed_buy_fee_sum BIGINT COMMENT 'Total jumlah pembayaran yang selesai pada hari ini', PRIMARY KEY (user_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Gunakan mekanisme penggabungan data agregasi untuk menghasilkan tabel agregat. 'fields.payed_buy_fee_sum.aggregate-function' = 'sum' -- Hitung jumlah data payed_buy_fee_sum untuk menghasilkan hasil agregat. -- Tabel dws_users tidak dikonsumsi oleh penyimpanan downstream dalam mode streaming. Oleh karena itu, Anda tidak perlu menentukan mekanisme pembuatan data inkremental. ); -- Buat tabel metrik agregat berdimensi pedagang. CREATE TABLE paimoncatalog.order_dw.dws_shops ( shop_id BIGINT, ds STRING, payed_buy_fee_sum BIGINT COMMENT 'Total jumlah pembayaran yang selesai pada hari ini', uv BIGINT COMMENT 'Jumlah total pengguna yang membeli komoditas pada hari ini', pv BIGINT COMMENT 'Jumlah total pembelian yang dilakukan oleh semua pengguna pada hari ini', PRIMARY KEY (shop_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Gunakan mekanisme penggabungan data agregasi untuk menghasilkan tabel agregat. 'fields.payed_buy_fee_sum.aggregate-function' = 'sum', -- Hitung jumlah data payed_buy_fee_sum untuk menghasilkan hasil agregat. 'fields.uv.aggregate-function' = 'sum', -- Hitung jumlah data unique visitor (UV) untuk menghasilkan hasil agregat. 'fields.pv.aggregate-function' = 'sum' -- Hitung jumlah data page view (PV) untuk menghasilkan hasil agregat. -- Tabel dws_shops tidak dikonsumsi oleh penyimpanan downstream dalam mode streaming. Oleh karena itu, Anda tidak perlu menentukan mekanisme pembuatan data inkremental. ); -- Untuk menghitung data pada tabel agregat dari perspektif pengguna dan data pada tabel agregat dari perspektif pedagang secara bersamaan, buat tabel antara yang menggunakan bidang user_id dan shop_id sebagai primary key. CREATE TABLE paimoncatalog.order_dw.dwm_users_shops ( user_id STRING, shop_id BIGINT, ds STRING, payed_buy_fee_sum BIGINT COMMENT 'Jumlah total yang dibayarkan pengguna di toko pada hari ini', pv BIGINT COMMENT 'Jumlah pembelian yang dilakukan pengguna di toko pada hari ini', PRIMARY KEY (user_id, shop_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Gunakan mekanisme penggabungan data agregasi untuk menghasilkan tabel agregat. 'fields.payed_buy_fee_sum.aggregate-function' = 'sum', -- Hitung jumlah data payed_buy_fee_sum untuk menghasilkan hasil agregat. 'fields.pv.aggregate-function' = 'sum', -- Hitung jumlah data PV untuk menghasilkan hasil agregat. 'changelog-producer' = 'lookup', -- Gunakan mekanisme pembuatan data inkremental lookup untuk menghasilkan changelog dengan latensi rendah. -- Umumnya, tabel antara pada lapisan DWM tidak menyediakan kueri untuk aplikasi lapisan atas. Oleh karena itu, performa penulisan dapat dioptimalkan. 'file.format' = 'avro', -- Gunakan format penyimpanan berorientasi baris Avro untuk memberikan performa penulisan yang lebih efisien. 'metadata.stats-mode' = 'none' -- Buang informasi statistik. Setelah informasi statistik dibuang, biaya kueri OLAP (online analytical processing) meningkat tetapi performa penulisan meningkat. Hal ini tidak berdampak pada pemrosesan aliran berkelanjutan. );Jika pesan
Query has been executeddikembalikan, berarti tabel berhasil dibuat.Konsumsi changelog dari tabel dwd_orders pada lapisan DWD.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada halaman yang muncul, klik untuk membuat draft bernama
dwm. Pada editor SQL, salin dan tempel kode berikut, lalu deploy draft tersebut. Mulai deployment pekerjaan dengan status awal.Pekerjaan ini menulis data dari tabel
dwd_orderske tabeldwm_users_shops. Dalam proses ini, mekanisme penggabungan data agregasi Apache Paimon digunakan untuk menghitung jumlah dataorder_feeguna memperoleh jumlah total konsumsi pengguna di toko. Pekerjaan ini juga menghitung jumlah entri data untuk memperoleh jumlah pembelian pengguna di toko.SET 'execution.checkpointing.max-concurrent-checkpoints' = '3'; SET 'table.exec.sink.upsert-materialize' = 'NONE'; SET 'execution.checkpointing.interval' = '10s'; SET 'execution.checkpointing.min-pause' = '10s'; INSERT INTO paimoncatalog.order_dw.dwm_users_shops SELECT order_user_id, order_shop_id, DATE_FORMAT (pay_create_time, 'yyyyMMdd') as ds, order_fee, 1 -- Satu catatan input mewakili satu pembelian. FROM paimoncatalog.order_dw.dwd_orders WHERE pay_id IS NOT NULL AND order_fee IS NOT NULL;Konsumsi changelog dari tabel dwm_users_shops pada lapisan DWM secara real-time.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada halaman yang muncul, klik untuk membuat draft bernama
dws. Pada editor SQL, salin dan tempel kode berikut, lalu deploy draft tersebut. Mulai deployment pekerjaan dengan status awal.Pekerjaan ini menulis data dari tabel
dwm_users_shopske tabeldws_usersdandws_shops. Dalam proses ini, mekanisme penggabungan data agregasi Apache Paimon digunakan untuk menghitung jumlah datapayed_buy_fee_sumpada tabeldws_usersguna memperoleh jumlah total konsumsi pengguna di semua toko, serta menghitung jumlah datapayed_buy_fee_sumpada tabeldws_shopsguna memperoleh jumlah total transaksi bisnis toko. Deployment ini juga menghitung jumlah entri data untuk memperoleh jumlah pengguna yang melakukan pembelian di toko, serta menghitung jumlah data PV untuk memperoleh jumlah total pembelian yang dilakukan semua pengguna di toko.SET 'execution.checkpointing.max-concurrent-checkpoints' = '3'; SET 'table.exec.sink.upsert-materialize' = 'NONE'; SET 'execution.checkpointing.interval' = '10s'; SET 'execution.checkpointing.min-pause' = '10s'; -- Berbeda dengan lapisan DWD, beberapa pernyataan INSERT yang menulis data ke tabel Apache Paimon yang berbeda dapat ditempatkan dalam deployment yang sama pada lapisan DWM. BEGIN STATEMENT SET; INSERT INTO paimoncatalog.order_dw.dws_users SELECT user_id, ds, payed_buy_fee_sum FROM paimoncatalog.order_dw.dwm_users_shops; -- Kolom shop_id digunakan sebagai primary key. Jumlah data toko populer tertentu mungkin jauh lebih tinggi daripada jumlah data toko lainnya. -- Oleh karena itu, merge lokal digunakan untuk mengagregasi data di memori sebelum data ditulis ke Apache Paimon. Hal ini membantu mengurangi masalah kesenjangan data. INSERT INTO paimoncatalog.order_dw.dws_shops /*+ OPTIONS('local-merge-buffer-size' = '64mb') */ SELECT shop_id, ds, payed_buy_fee_sum, 1, -- Satu catatan input mewakili seluruh konsumsi pengguna di toko. pv FROM paimoncatalog.order_dw.dwm_users_shops; END;Lihat data tabel dws_users dan dws_shops.
Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada tab Scripts, klik . Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu Run:
-- Lihat data tabel dws_users. SELECT * FROM paimoncatalog.order_dw.dws_users ORDER BY user_id;
-- Lihat data tabel dws_shops. SELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;
Tangkap perubahan di database bisnis
Danau data terpadu streaming telah dibuat. Bagian ini menguji kemampuan danau data terpadu streaming untuk menangkap perubahan di database bisnis.
Masukkan data berikut ke database
order_dwMySQL:INSERT INTO orders VALUES (100008, 'user_001', 12345, 3, 3000, '2023-02-15 17:40:56', '2023-02-15 18:42:56', 1), (100009, 'user_002', 12348, 4, 1000, '2023-02-15 18:40:56', '2023-02-15 19:42:56', 1), (100010, 'user_003', 12348, 2, 2000, '2023-02-15 19:40:56', '2023-02-15 20:42:56', 1); INSERT INTO orders_pay VALUES (2008, 100008, 1, '2023-02-15 18:40:56'), (2009, 100009, 1, '2023-02-15 19:40:56'), (2010, 100010, 0, '2023-02-15 20:40:56');Lihat data tabel dws_users dan dws_shops.Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih . Pada tab Scripts, klik . Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu Run:
Tabel
dws_usersSELECT * FROM paimoncatalog.order_dw.dws_users ORDER BY user_id;
Tabel
dws_shopsSELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;
Gunakan danau data terpadu streaming
Bagian sebelumnya menjelaskan cara membuat katalog Apache Paimon dan menulis data ke tabel Apache Paimon di Konsol Realtime Compute for Apache Flink. Bagian ini menjelaskan skenario sederhana spesifik di mana StarRocks digunakan untuk menganalisis data setelah danau data terpadu streaming dibuat.
Login ke instans StarRocks dan buat katalog Apache Paimon.
CREATE EXTERNAL CATALOG paimon_catalog
PROPERTIES
(
'type' = 'paimon',
'paimon.catalog.type' = 'filesystem',
'aliyun.oss.endpoint' = 'oss-cn-beijing-internal.aliyuncs.com',
'paimon.catalog.warehouse' = 'oss://<bucket>/<object>'
);Parameter | Wajib | Deskripsi |
type | Ya | Jenis sumber data. Atur parameter ke paimon. |
paimon.catalog.type | Ya | Jenis penyimpanan metadata yang digunakan oleh katalog Apache Paimon. Pada contoh ini, digunakan filesystem. |
aliyun.oss.endpoint | Ya | Titik akhir OSS atau OSS-HDFS. Parameter ini wajib jika Anda mengatur parameter paimon.catalog.warehouse ke path OSS atau OSS-HDFS. |
paimon.catalog.warehouse | Ya | Direktori gudang data yang ditentukan di OSS. Formatnya adalah oss://<bucket>/<object>. Di mana:
Anda dapat melihat nama bucket dan nama objek di Konsol OSS. |
Lakukan kueri peringkat
Analisis tabel agregat pada lapisan DWS. Kode contoh berikut menunjukkan cara menggunakan StarRocks untuk mengkueri tiga toko teratas dengan volume transaksi tertinggi pada 15 Februari 2023:
SELECT ROW_NUMBER() OVER (ORDER BY payed_buy_fee_sum DESC) AS rn, shop_id, payed_buy_fee_sum
FROM dws_shops
WHERE ds = '20230215'
ORDER BY rn LIMIT 3;
Lakukan kueri detail
Analisis tabel lebar pada lapisan DWD. Kode contoh berikut menunjukkan cara menggunakan StarRocks untuk mengkueri detail pesanan yang dibayar pelanggan pada platform pembayaran tertentu pada Februari 2023:
SELECT * FROM dwd_orders
WHERE order_create_time >= '2023-02-01 00:00:00' AND order_create_time < '2023-03-01 00:00:00'
AND order_user_id = 'user_001'
AND pay_platform = 0
ORDER BY order_create_time;;
Kueri laporan data
Analisis tabel lebar pada lapisan DWD. Kode contoh berikut menunjukkan cara menggunakan StarRocks untuk mengkueri jumlah total pesanan dan jumlah total pesanan untuk setiap kategori pada Februari 2023:
SELECT
order_create_time AS order_create_date,
order_product_catalog_name,
COUNT(*),
SUM(order_fee)
FROM
dwd_orders
WHERE
order_create_time >= '2023-02-01 00:00:00' and order_create_time < '2023-03-01 00:00:00'
GROUP BY
order_create_date, order_product_catalog_name
ORDER BY
order_create_date, order_product_catalog_name;