Pelajari cara membangun danau data terpadu streaming menggunakan Realtime Compute for Apache Flink, Apache Paimon, dan StarRocks.
Latar Belakang
Seiring meningkatnya digitalisasi bisnis, permintaan akan data yang tepat waktu terus bertambah. Metode tradisional dalam membangun gudang data offline melibatkan penjadwalan pekerjaan offline untuk secara berkala menggabungkan data baru ke dalam struktur berlapis, seperti ODS, DWD, DWS, dan ADS. Namun, pendekatan ini memiliki dua masalah signifikan: latensi tinggi dan biaya tinggi. Pekerjaan offline biasanya dijadwalkan per jam atau bahkan harian, sehingga konsumen downstream hanya dapat melihat data yang berumur minimal satu jam atau satu hari. Selain itu, pembaruan sering kali memerlukan penimpaan seluruh partisi. Proses yang tidak efisien ini membaca ulang semua data asli dalam suatu partisi untuk menggabungkan perubahan baru.
Membangun danau data terpadu streaming dengan Realtime Compute for Apache Flink dan Apache Paimon mengatasi masalah latensi dan biaya pada gudang data offline tradisional. Kemampuan pemrosesan real-time Flink memungkinkan data mengalir secara terus-menerus antar lapisan gudang. Sementara itu, kemampuan pembaruan efisien Paimon mengirimkan perubahan data ke konsumen downstream hanya dengan latensi tingkat menit. Akibatnya, danau data terpadu streaming ini menawarkan keunggulan signifikan baik dalam hal latensi rendah maupun efisiensi biaya.
Untuk informasi lebih lanjut tentang Apache Paimon, lihat Fitur dan situs web resmi Apache Paimon.
Arsitektur dan manfaat
Arsitektur
Realtime Compute for Apache Flink adalah mesin pemrosesan aliran yang andal untuk memproses volume besar data real-time secara efisien. Paimon adalah format penyimpanan danau terpadu untuk pemrosesan aliran dan batch yang mendukung pembaruan throughput tinggi serta kueri latensi rendah. Paimon terintegrasi erat dengan Flink untuk menyediakan solusi danau data terpadu streaming all-in-one. Berikut ini menunjukkan arsitektur pembangunan danau data terpadu streaming menggunakan Flink dan Paimon.
-
Flink menulis data dari sumber data ke Paimon, membuat lapisan ODS.
-
Flink berlangganan changelog dari lapisan ODS, mentransformasi data, lalu menulis kembali data tersebut ke Paimon sebagai lapisan DWD.
-
Flink berlangganan changelog dari lapisan DWD, mentransformasi data, lalu menulis kembali data tersebut ke Paimon sebagai lapisan DWS.
-
Akhirnya, StarRocks di E-MapReduce membaca tabel eksternal Paimon untuk melayani kueri aplikasi.

Manfaat
Solusi ini menawarkan manfaat berikut:
-
Setiap lapisan data di Paimon menyebarkan perubahan ke sistem downstream dalam hitungan menit. Hal ini mengurangi latensi gudang data offline tradisional dari hitungan jam atau bahkan hari menjadi hitungan menit.
-
Setiap lapisan data di Paimon langsung mengingesti data perubahan tanpa menimpa partisi. Ini secara signifikan mengurangi biaya pembaruan dan koreksi data pada gudang data offline tradisional serta menyelesaikan tantangan dalam mengkueri, memperbarui, dan mengoreksi data di lapisan antara.
-
Solusi ini memiliki model terpadu dan arsitektur yang disederhanakan. Logika pipeline ETL diimplementasikan menggunakan Flink SQL. Data di lapisan ODS, DWD, dan DWS disimpan dalam format Paimon. Hal ini mengurangi kompleksitas arsitektur dan meningkatkan efisiensi pemrosesan data.
Solusi ini bergantung pada tiga kemampuan inti Paimon, sebagaimana dirinci dalam tabel berikut.
|
Kemampuan inti |
Deskripsi |
|
Primary key table updates |
Paimon secara internal menggunakan struktur data LSM tree untuk memungkinkan pembaruan data yang efisien. Untuk informasi lebih lanjut tentang tabel primary key Paimon dan struktur data dasarnya, lihat Primary Key Table dan File Layouts. |
|
Changelog Producer |
Paimon dapat menghasilkan data inkremental lengkap untuk setiap aliran data input, di mana setiap catatan |
|
Merge Engine |
Saat tabel primary key Paimon menerima beberapa catatan dengan primary key yang sama, Merge Engine menggabungkan catatan-catatan tersebut menjadi satu catatan untuk menjaga keunikan kunci. Paimon mendukung berbagai perilaku penggabungan, seperti |
Kasus penggunaan
Artikel ini menggunakan platform e-commerce sebagai contoh untuk menunjukkan cara membangun danau data terpadu streaming. Solusi ini memproses dan membersihkan data untuk mendukung kueri dari aplikasi upstream. Dengan menggunakan pelapisan data, solusi ini memungkinkan penggunaan ulang data di berbagai skenario bisnis, seperti Dasbor transaksi, analisis data perilaku, tag profil pengguna, dan rekomendasi personalisasi.

-
Bangun lapisan ODS: Ingesti real-time dari database bisnis
Flink mengingesti tiga tabel bisnis dari MySQL, yaituorders,orders_pay, danproduct_catalog, ke OSS secara real time. Data ini, yang disimpan dalam format Paimon, membentuk lapisan ODS. -
Bangun lapisan DWD: Tabel lebar bertema
Tabelorders,product_catalog, danorders_paydigabungkan menggunakan mekanisme merge partial-update Paimon, menghasilkan tabel lebar bertema untuk lapisan DWD dan changelog dengan latensi tingkat menit. -
Bangun lapisan DWS: Perhitungan metrik
Flink mengonsumsi changelog dari tabel lebar secara real time. Flink menggunakan mekanisme merge pre-aggregation Paimon untuk membuat tabel agregasi antara,dwm_users_shops. Proses ini akhirnya menghasilkan dua tabel untuk lapisan DWS:dws_usersuntuk metrik agregasi pengguna dandws_shopsuntuk metrik agregasi toko.
Prasyarat
-
Aktifkan Data Lake Formation (DLF). Kami merekomendasikan penggunaan DLF 2.5 sebagai layanan penyimpanan. Untuk informasi lebih lanjut, lihat Memulai dengan DLF.
-
Buat ruang kerja Realtime Compute for Apache Flink. Untuk informasi lebih lanjut, lihat Buat ruang kerja.
-
Buat instans Serverless StarRocks. Untuk informasi lebih lanjut, lihat Memulai dengan instans Serverless StarRocks.
Instans StarRocks, DLF, dan ruang kerja Flink Anda harus berada di Wilayah yang sama.
Batasan
Solusi danau data terpadu streaming ini memerlukan Ververica Runtime (VVR) 11.1.0 atau yang lebih baru.
Bangun danau data terpadu streaming
Siapkan sumber data MySQL CDC
Tutorial ini menggunakan ApsaraDB RDS for MySQL sebagai contoh. Anda membuat database bernama order_dw dan tiga tabel dengan data sampel.
-
(Usang. Dialihkan ke Langkah 1.) Buat instans ApsaraDB RDS for MySQL.
PentingPastikan instans ApsaraDB RDS for MySQL dan ruang kerja Realtime Compute for Apache Flink Anda berada di VPC yang sama. Jika berada di VPC berbeda, lihat Bagaimana cara mengakses layanan lain lintas VPC?
-
(Usang. Dialihkan ke Langkah 1.) Buat database dan akun.
Buat database bernama order_dw. Lalu, buat akun istimewa atau akun standar dengan izin baca dan tulis untuk database order_dw.
Buat tiga tabel dan masukkan data ke dalamnya.
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 ); -- Prepare 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');
Manage metadata
Buat katalog Paimon
-
Login ke Konsol Realtime Compute for Apache Flink.
-
Di panel navigasi kiri, buka Metadata Management dan klik Create Catalog.
-
Di tab Built-in Catalog, klik Apache Paimon, lalu klik Next.
-
Konfigurasikan parameter berikut, pilih DLF sebagai tipe penyimpanan, lalu klik OK.
Parameter
Deskripsi
Wajib
Catatan
metastore
Jenis metastore.
Ya
Dalam contoh ini, pilih DLF.
catalog name
Nama katalog data DLF.
PentingJika Anda menggunakan Pengguna RAM atau Peran RAM, pastikan Pengguna RAM atau Peran RAM tersebut memiliki izin yang diperlukan untuk membaca dan menulis data ke DLF. Untuk informasi lebih lanjut, lihat Mengelola izin data.
Ya
Kami merekomendasikan penggunaan DLF 2.5. Versi ini menghilangkan kebutuhan untuk memasukkan informasi seperti Pasangan Kunci Akses dan memungkinkan Anda memilih katalog data DLF yang sudah ada dengan cepat. Untuk informasi cara membuat katalog data, lihat Katalog data.
Setelah membuat katalog data bernama paimoncatalog, pilih dari daftar.
-
Di katalog data, buat database bernama order_dw untuk menyinkronkan semua tabel dari database MySQL order_dw.
Di panel navigasi kiri, pilih , lalu klik New untuk membuka jendela kueri sementara.
-- Gunakan sumber data paimoncatalog USE CATALOG paimoncatalog; -- Buat database order_dw CREATE DATABASE order_dw;Pesan yang dikembalikan
The following statement has been executed successfully!menunjukkan bahwa database berhasil dibuat.
Untuk informasi lebih lanjut tentang cara menggunakan katalog Paimon, lihat Mengelola katalog Paimon.
Buat katalog MySQL
-
Di halaman Metadata Management, klik Create Catalog.
-
Di tab Built-in Catalog, klik MySQL, lalu klik Next.
-
Konfigurasikan parameter berikut dan klik OK untuk membuat katalog MySQL bernama mysqlcatalog.
Parameter
Deskripsi
Wajib
Catatan
catalog name
Nama katalog.
Ya
Masukkan nama kustom. Tutorial ini menggunakan mysqlcatalog.
hostname
Alamat IP atau hostname database MySQL.
Ya
Untuk informasi lebih lanjut, lihat Melihat dan mengelola titik akhir dan port koneksi instans. Karena instans ApsaraDB RDS for MySQL dan ruang kerja Realtime Compute for Apache Flink berada di VPC yang sama, masukkan Titik akhir internal.
port
Nomor port layanan database MySQL. Nilai default adalah 3306.
Tidak
Untuk informasi lebih lanjut, lihat Melihat dan mengelola titik akhir dan port koneksi instans.
default-database
Nama database MySQL default.
Ya
Masukkan nama database yang akan disinkronkan, yaitu order_dw dalam tutorial ini.
username
Username untuk layanan database MySQL.
Ya
Gunakan akun yang dibuat di Menyiapkan sumber data MySQL CDC.
password
Password untuk layanan database MySQL.
Ya
Gunakan password untuk akun yang dibuat di Menyiapkan sumber data MySQL CDC.
Bangun lapisan ODS: Ingesti real-time
Bangun lapisan ODS dalam satu langkah menggunakan pekerjaan ingesti data Flink CDC dalam format YAML untuk menyinkronkan data dari MySQL ke Paimon.
-
Buat dan mulai pekerjaan ingesti data YAML.
-
Di Konsol Realtime Compute for Apache Flink, buka halaman dan buat pekerjaan draft YAML kosong bernama ods.
-
Salin kode berikut ke editor dan perbarui parameter seperti username dan password.
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.\.* # Gunakan ekspresi reguler untuk membaca semua tabel di database order_dw. # (Opsional) Sinkronkan data dari tabel yang baru dibuat selama fase inkremental. scan.binlog.newly-added-table.enabled: true # (Opsional) Sinkronkan komentar tabel dan bidang. include-comments.enabled: true # (Opsional) Utamakan distribusi chunk tak terbatas untuk mencegah potensi masalah OutOfMemory TaskManager. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Opsional) Aktifkan filter parsing untuk mempercepat pembacaan. scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: paimon name: Paimon Sink catalog.properties.metastore: rest catalog.properties.uri: http://cn-beijing-vpc.dlf.aliyuncs.com catalog.properties.warehouse: paimoncatalog catalog.properties.token.provider: dlf pipeline: name: MySQL to Paimon PipelineParameter
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 akses untuk DLF Rest Catalog Server adalah
http://[region-id]-vpc.dlf.aliyuncs.com. Untuk informasi lebih lanjut tentang Region ID, lihat Titik Akhir Layanan.Ya
http://cn-beijing-vpc.dlf.aliyuncs.com
catalog.properties.warehouseNama katalog DLF.
Ya
paimoncatalog
hostnameAlamat IP atau hostname database MySQL. Untuk informasi lebih lanjut, lihat Melihat dan mengelola titik akhir dan port koneksi instans. Karena instans ApsaraDB RDS for MySQL dan ruang kerja Realtime Compute for Apache Flink berada di VPC yang sama, masukkan Titik akhir internal.
Ya
rm-bp1e********566g.mysql.rds.aliyuncs.com
usernameUsername untuk database MySQL. Kami merekomendasikan penggunaan manajemen rahasia. Untuk informasi lebih lanjut, lihat Mengelola variabel.
Ya
${secret_values.username}
passwordPassword untuk database MySQL. Kami merekomendasikan penggunaan manajemen rahasia. Untuk informasi lebih lanjut, lihat Mengelola variabel.
Ya
${secret_values.password}
Untuk mengoptimalkan kinerja penulisan Paimon, lihat Pengaturan kinerja Paimon.
-
Di pojok kanan atas, klik Deploy.
-
Di , klik Start di kolom Actions pekerjaan ODS yang baru dideploy, lalu pilih stateless start untuk memulai pekerjaan. Untuk informasi lebih lanjut tentang konfigurasi startup pekerjaan, lihat Job Startup.
-
-
Lihat data di tiga tabel yang telah disinkronkan dari MySQL ke Paimon.
Di Konsol Realtime Compute for Apache Flink, buka halaman . Di tab Query Editor, salin kode berikut ke editor, pilih pernyataan tersebut, lalu klik Run di pojok kanan atas.
SELECT * FROM paimoncatalog.order_dw.orders ORDER BY order_id;Setelah kueri dieksekusi, hasilnya mengembalikan tujuh catatan pesanan dari tabel orders. Hasil mencakup kolom order_id (100001 hingga 100007), user_id, shop_id, product_id, buy_fee (1000 hingga 5000), create_time, update_time, dan state. Nilai state adalah 1 untuk semua catatan.
Bangun lapisan DWD: tabel lebar bertema
-
Buat tabel lebar dwd_orders
Di Konsol Realtime Compute for Apache Flink, buka halaman . Di tab Query Editor, salin kode berikut ke editor, pilih pernyataan tersebut, lalu klik Run di pojok kanan atas.
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 mesin merge partial-update untuk menghasilkan tabel lebar 'changelog-producer' = 'lookup' -- Gunakan produsen changelog lookup untuk menghasilkan changelog dengan latensi rendah );Pesan
The following statement has been executed successfully!menunjukkan pembuatan berhasil. -
Konsumsi data perubahan ODS
Di Konsol Realtime Compute for Apache Flink, buka halaman . Buat pekerjaan streaming SQL baru bernama dwd, salin kode berikut ke editor SQL, klik Deploy, lalu mulai pekerjaan menggunakan Stateless Start.
Pekerjaan SQL ini melakukan join tabel orders dengan tabel product_catalog dalam join tabel dimensi. Hasilnya kemudian ditulis ke tabel dwd_orders bersama dengan data dari tabel orders_pay. Mekanisme merge partial-update Paimon melakukan join catatan dari tabel orders dan orders_pay yang memiliki order_id yang sama.
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'; -- Paimon saat ini tidak mendukung beberapa pernyataan INSERT ke tabel yang sama dalam satu pekerjaan. Oleh karena itu, gunakan UNION ALL di sini. 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 di dwd_orders
Di Konsol Realtime Compute for Apache Flink, buka halaman . Di tab Query Editor, salin kode berikut ke editor, pilih pernyataan tersebut, lalu klik Run di pojok kanan atas.
SELECT * FROM paimoncatalog.order_dw.dwd_orders ORDER BY order_id;Setelah kueri berhasil dieksekusi, hasilnya mengembalikan tujuh catatan pesanan dari tabel lebar
dwd_orders. Hasil mencakup delapan bidang:order_id,order_user_id,order_shop_id,order_product_id,order_product_catalog_name,order_fee,order_create_time, danorder_update_time.order_idberkisar dari 100001 hingga 100007, biaya pesanan berkisar dari 1000 hingga 5000, dan semua waktu pembuatan berada pada tanggal 2023-02-15.
Bangun lapisan DWS: perhitungan metrik
-
Buat tabel agregat lapisan DWS dws_users dan dws_shops.
Di Konsol Realtime Compute for Apache Flink, buka halaman . Di tab Query Editor, salin kode berikut ke editor, pilih pernyataan tersebut, lalu klik Run di pojok kanan atas.
-- Tabel agregat dimensi pengguna. CREATE TABLE paimoncatalog.order_dw.dws_users ( user_id STRING, ds STRING, paid_buy_fee_sum BIGINT COMMENT 'Total jumlah yang dibayarkan pada hari ini', PRIMARY KEY (user_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Gunakan mesin merge agregasi untuk menghasilkan tabel agregat. 'fields.paid_buy_fee_sum.aggregate-function' = 'sum' -- Agregasi hasil dengan menjumlahkan data di paid_buy_fee_sum. -- Karena tabel dws_users tidak dikonsumsi oleh pekerjaan streaming downstream, Anda tidak perlu menentukan produsen changelog. ); -- Tabel agregat dimensi toko. CREATE TABLE paimoncatalog.order_dw.dws_shops ( shop_id BIGINT, ds STRING, paid_buy_fee_sum BIGINT COMMENT 'Total jumlah yang dibayarkan pada hari ini', uv BIGINT COMMENT 'Jumlah total pengguna pembelian unik pada hari ini', pv BIGINT COMMENT 'Jumlah total pembelian oleh pengguna pada hari ini', PRIMARY KEY (shop_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Gunakan mesin merge agregasi untuk menghasilkan tabel agregat. 'fields.paid_buy_fee_sum.aggregate-function' = 'sum', -- Agregasi hasil dengan menjumlahkan data di paid_buy_fee_sum. 'fields.uv.aggregate-function' = 'sum', -- Agregasi hasil dengan menjumlahkan data di uv. 'fields.pv.aggregate-function' = 'sum' -- Agregasi hasil dengan menjumlahkan data di pv. -- Karena tabel dws_shops tidak dikonsumsi oleh pekerjaan streaming downstream, Anda tidak perlu menentukan produsen changelog. ); -- Untuk menghitung tabel agregat dimensi pengguna dan toko secara simultan, buat tabel antara dengan primary key pengguna dan toko. CREATE TABLE paimoncatalog.order_dw.dwm_users_shops ( user_id STRING, shop_id BIGINT, ds STRING, paid_buy_fee_sum BIGINT COMMENT 'Total jumlah yang dibayarkan pengguna di toko pada hari ini', pv BIGINT COMMENT 'Jumlah kali pengguna membeli dari toko pada hari ini', PRIMARY KEY (user_id, shop_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Gunakan mesin merge agregasi untuk menghasilkan tabel agregat. 'fields.paid_buy_fee_sum.aggregate-function' = 'sum', -- Agregasi hasil dengan menjumlahkan data di paid_buy_fee_sum. 'fields.pv.aggregate-function' = 'sum', -- Agregasi hasil dengan menjumlahkan data di pv. 'changelog-producer' = 'lookup', -- Gunakan produsen changelog lookup untuk menghasilkan changelog dengan latensi rendah. -- Tabel antara di lapisan DWM biasanya tidak dikueri langsung oleh aplikasi upstream, sehingga dapat dioptimalkan untuk kinerja penulisan. 'file.format' = 'avro', -- Menggunakan format penyimpanan berbasis baris Avro memberikan kinerja penulisan yang lebih efisien. 'metadata.stats-mode' = 'none' -- Menghilangkan statistik meningkatkan biaya kueri OLAP (tanpa dampak pada pemrosesan aliran kontinu) tetapi meningkatkan kinerja penulisan. );Pesan
The following statement has been executed successfully!menunjukkan pembuatan berhasil. -
Konsumsi data perubahan dari tabel lapisan DWD dwd_orders.
Di Konsol Realtime Compute for Apache Flink, buka halaman , buat pekerjaan streaming SQL baru bernama dwm. Salin kode berikut ke editor SQL, klik Deploy, lalu mulai pekerjaan menggunakan Stateless Start.
Pekerjaan SQL ini menulis data dari tabel dwd_orders ke tabel dwm_users_shops. Mekanisme merge data pre-aggregation Paimon secara otomatis menjumlahkan order_fee untuk menghitung total jumlah yang dihabiskan pengguna di toko. Paimon juga menjumlahkan nilai 1 untuk menghitung jumlah kali pengguna melakukan pembelian di toko tersebut.
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 merepresentasikan satu pembelian FROM paimoncatalog.order_dw.dwd_orders WHERE pay_id IS NOT NULL AND order_fee IS NOT NULL; -
Konsumsi data perubahan dari tabel lapisan DWM dwm_users_shops secara real time.
Di Konsol Realtime Compute for Apache Flink, buka halaman . Buat pekerjaan streaming SQL baru bernama dws. Salin kode berikut ke editor SQL, klik Deploy, lalu mulai pekerjaan menggunakan Stateless Start.
Pekerjaan SQL ini menulis data dari tabel dwm_users_shops ke tabel dws_users dan dws_shops. Mekanisme merge data pre-aggregation Paimon digunakan untuk menghitung total pengeluaran setiap pengguna (paid_buy_fee_sum) di tabel dws_users. Di tabel dws_shops, Paimon menghitung total penjualan toko (paid_buy_fee_sum), jumlah pengguna pembeli (dengan menjumlahkan 1), dan jumlah total pembelian (pv).
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 pekerjaan DWD, setiap pernyataan INSERT di sini menulis ke tabel Paimon yang berbeda, sehingga dapat dimasukkan dalam pekerjaan yang sama. BEGIN STATEMENT SET; INSERT INTO paimoncatalog.order_dw.dws_users SELECT user_id, ds, paid_buy_fee_sum FROM paimoncatalog.order_dw.dwm_users_shops; -- Primary key-nya adalah ID toko. Data untuk toko populer mungkin jauh lebih besar daripada toko lainnya. -- Oleh karena itu, gunakan local merge untuk melakukan pre-aggregasi data di memori sebelum menulis ke Paimon guna mengurangi kesenjangan data. INSERT INTO paimoncatalog.order_dw.dws_shops /*+ OPTIONS('local-merge-buffer-size' = '64mb') */ SELECT shop_id, ds, paid_buy_fee_sum, 1, -- Satu catatan input merepresentasikan semua pembelian pengguna di toko ini pv FROM paimoncatalog.order_dw.dwm_users_shops; END; -
Lihat data di tabel dws_users dan dws_shops.
Di Konsol Realtime Compute for Apache Flink, buka halaman . Di tab Query Editor, salin kode berikut ke editor, pilih pernyataan tersebut, lalu klik Run di pojok kanan atas.
--Lihat data di tabel dws_users SELECT * FROM paimoncatalog.order_dw.dws_users ORDER BY user_id;Setelah kueri dieksekusi, tabel dws_users mengembalikan tiga catatan dengan kolom berikut:
user_id,ds, danpaid_buy_fee_sum. Catatan tersebut adalah:user_001/ 20230215 / 8000,user_002/ 20230215 / 5000, danuser_003/ 20230215 / 5000.--Lihat data di tabel dws_shops SELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;Setelah Anda mengkueri tabel dws_shops, hasilnya berisi empat baris dan lima kolom:
shop_id,ds,paid_buy_fee_sum,uv, danpv. shop_id adalah 12345 hingga 12348, ds adalah 20230215 untuk semua catatan, nilai paid_buy_fee_sum adalah 5000, 4000, 7000, dan 2000, nilai uv adalah 1, 1, 2, dan 2, serta nilai pv adalah 1, 1, 3, dan 2, masing-masing.
Tangkap perubahan dari database bisnis
Setelah membangun danau data terpadu streaming, uji kemampuannya untuk menangkap perubahan dari database bisnis.
-
Masukkan data berikut ke database order_dw di MySQL.
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 di tabel dws_users dan dws_shops. Di Konsol Realtime Compute for Apache Flink, buka halaman . Di tab Query Editor, salin kode berikut ke editor, pilih pernyataan tersebut, lalu klik Run di pojok kanan atas.
-
tabel dws_users
SELECT * FROM paimoncatalog.order_dw.dws_users ORDER BY user_id;Kueri mengembalikan tiga baris dengan kolom berikut:
user_id,ds, danpaid_buy_fee_sum. Catatan tersebut adalah:user_001/ 20230215 / 11000,user_002/ 20230215 / 6000, danuser_003/ 20230215 / 7000. -
tabel dws_shops
SELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;Kueri mengembalikan empat baris dengan lima kolom:
shop_id,ds,paid_buy_fee_sum,uv, danpv. shop_id adalah 12345, 12346, 12347, dan 12348, ds adalah 20230215 untuk semua catatan, nilai paid_buy_fee_sum adalah 8000, 4000, 7000, dan 5000, nilai uv adalah 1, 1, 2, dan 3, serta nilai pv adalah 2, 1, 3, dan 4, masing-masing.
-
Gunakan danau data terpadu streaming
Bagian sebelumnya menunjukkan cara membuat katalog Paimon dan menulis ke tabel Paimon di Flink. Bagian ini menunjukkan beberapa kasus penggunaan analisis data sederhana menggunakan StarRocks setelah Anda menyiapkan danau data terpadu streaming.
Hubungkan StarRocks ke DLF
Untuk informasi lebih lanjut, lihat Mengakses DLF dari Serverless StarRocks.
Kueri peringkat
Kueri StarRocks berikut menganalisis tabel agregat di lapisan DWS untuk mengambil tiga toko teratas berdasarkan jumlah transaksi pada 15 Februari 2023.
SELECT ROW_NUMBER() OVER (ORDER BY paid_buy_fee_sum DESC) AS rn, shop_id, paid_buy_fee_sum
FROM paimoncatalog.order_dw.dws_shops
WHERE ds = '20230215'
ORDER BY rn LIMIT 3;
Kueri mengembalikan tiga baris: rn=1 sesuai dengan shop_id=12345 dan paid_buy_fee_sum=8000.00; rn=2 sesuai dengan shop_id=12347 dan paid_buy_fee_sum=7000.00; dan rn=3 sesuai dengan shop_id=12348 dan paid_buy_fee_sum=5000.00.
Detail kueri
Kueri StarRocks berikut menganalisis tabel lebar di lapisan DWD untuk mengambil detail pesanan pelanggan tertentu yang menggunakan platform pembayaran tertentu pada Februari 2023.
SELECT * FROM paimoncatalog.order_dw.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 mengembalikan hasil berikut.
order_id order_user_id order_shop_id order_product_id order_product_catalog_name order_fee
0 100006 user_001 12348 1 phone_aaa 1000
1 100004 user_001 12347 4 phone_ddd 2000
Laporan data
Kueri StarRocks berikut menganalisis tabel lebar di lapisan DWD untuk menghasilkan laporan jumlah total pesanan dan jumlah total pesanan untuk setiap tanggal dan kategori produk pada Februari 2023.
SELECT
order_create_time AS order_create_date,
order_product_catalog_name,
COUNT(*),
SUM(order_fee)
FROM
paimoncatalog.order_dw.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;
Kueri mengembalikan sembilan catatan, menunjukkan data pesanan per jam dari pukul 10:40:56 hingga 18:40:56 pada 15 Februari 2023. Kategori produk meliputi phone_ddd, phone_aaa, phone_eee, phone_ccc, dan phone_bbb. Setiap catatan memiliki count(*) sebesar 1, dan nilai sum(order_fee) masing-masing adalah 2000, 1000, 1000, 2000, 3000, 4000, 5000, 3000, dan 1000.
Referensi
-
Untuk membangun danau data offline Paimon menggunakan pemrosesan batch Flink, lihat Pemrosesan batch Flink.