All Products
Search
Document Center

Data Lake Formation:Bangun danau data terpadu streaming dengan DLF dan Flink

Last Updated:Jan 16, 2026

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.

  1. Flink menulis data dari sumber data ke Paimon untuk membentuk lapisan ODS.

  2. Flink berlangganan changelog data dari lapisan ODS, memprosesnya, lalu menulis kembali ke Paimon untuk membentuk lapisan DWD.

  3. Flink berlangganan changelog data dari lapisan DWD, memprosesnya, lalu menulis kembali ke Paimon untuk membentuk lapisan DWS.

  4. Terakhir, StarRocks di E-MapReduce open source membaca tabel eksternal Paimon untuk menyediakan layanan kueri bagi aplikasi.

image

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.

image

  1. 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.

  2. 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.

  3. 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

Catatan

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.

  1. Buat instans ApsaraDB RDS for MySQL.

    Catatan

    Jika ApsaraDB RDS for MySQL dan ruang kerja Flink Anda tidak berada dalam VPC yang sama, lihat Bagaimana cara mengakses layanan lain lintas VPC?.

  2. Buat database dan akun.

    Buat database bernama order_dw dan buat akun istimewa atau akun standar yang memiliki izin baca dan tulis pada database order_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

  1. Login ke Konsol Realtime Compute for Apache Flink.

  2. Pada bilah navigasi di sebelah kiri, pilih Metadata Management dan klik Create Catalog.

  3. Pada tab Built-in Catalog, klik Apache Paimon, lalu klik Next.

  4. 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.

    Penting

    Jika 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 paimoncatalog telah dibuat sebelumnya.

  5. 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 Data Query > Query Script, 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

  1. Pada halaman Metadata Management, klik Create Catalog.

  2. Pada tab Built-in Catalog, klik MySQL, lalu klik Next.

  3. 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.

  1. Buat pekerjaan ingest data.

    1. Login ke Konsol Manajemen Realtime Compute for Apache Flink. Klik Console pada kolom Actions ruang kerja Anda untuk masuk ke Konsol Pengembangan.

    2. Pada menu navigasi kiri, pilih Development > ETL. Klik + > New Blank Stream Draft. Pada dialog New Draft, masukkan ods pada 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 Pipeline

      Item konfigurasi

      Deskripsi

      Wajib?

      Contoh

      catalog.properties.metastore

      Jenis metastore. Atur ke rest.

      Ya

      rest

      catalog.properties.token.provider

      Penyedia token. Atur ke dlf.

      Ya

      dlf

      catalog.properties.uri

      URI 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.warehouse

      Nama katalog DLF.

      Ya

      paimoncatalog

      Anda dapat mengonfigurasi properti tabel Paimon untuk meningkatkan performa penulisan. Untuk detailnya, lihat Optimalisasi performa.

    3. Pada pojok kanan atas editor SQL, klik Deploy.

    4. Pada menu navigasi kiri, pilih O&M > Deployments. Pada halaman Deployments, temukan deployment pekerjaan ods lalu klik Start pada kolom Actions. Pada panel Start Job, pilih Initial Mode lalu klik Start.

  2. Lihat data yang disinkronkan dari MySQL ke Paimon.

    Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih Development > Scripts. Klik + > New Script . Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu klik Run:

    SELECT * FROM paimoncatalog.order_dw.orders ORDER BY order_id;

    截屏2024-09-02 14

Bangun lapisan DWD: Tabel lebar

  1. Buat tabel lebar bernama dwd_orders pada lapisan DWD di Apache Paimon.

    Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih Development > Scripts. Pada tab Scripts, klik + > New Script. 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 executed dikembalikan, berarti tabel berhasil dibuat.

  2. 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 Development > ETL. Pada halaman yang muncul, klik + > New Blank Stream Draft 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 orders dengan tabel dimensi bernama product_catalog dan menulis hasil join serta data dari tabel orders_pay ke tabel lebar bernama dwd_orders. Dalam proses ini, mekanisme penggabungan data pembaruan parsial Apache Paimon digunakan untuk menggabungkan data yang memiliki nilai order_id yang sama pada tabel orders dan orders_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;
  3. Lihat data tabel lebar bernama dwd_orders.

    Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih Development > Scripts. Pada tab Scripts, klik + > New Script. Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu Run:

    SELECT * FROM paimoncatalog.order_dw.dwd_orders ORDER BY order_id;

    截屏2024-09-02 14

Bangun lapisan DWS: Perhitungan metrik

  1. 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 Development > Scripts. Pada tab Scripts, klik + > New Script. 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 executed dikembalikan, berarti tabel berhasil dibuat.

  2. Konsumsi changelog dari tabel dwd_orders pada lapisan DWD.

    Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih Development > ETL. Pada halaman yang muncul, klik + > New Blank Stream Draft 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_orders ke tabel dwm_users_shops. Dalam proses ini, mekanisme penggabungan data agregasi Apache Paimon digunakan untuk menghitung jumlah data order_fee guna 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;
  3. 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 Development > ETL. Pada halaman yang muncul, klik + > New Blank Stream Draft 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_shops ke tabel dws_users dan dws_shops. Dalam proses ini, mekanisme penggabungan data agregasi Apache Paimon digunakan untuk menghitung jumlah data payed_buy_fee_sum pada tabel dws_users guna memperoleh jumlah total konsumsi pengguna di semua toko, serta menghitung jumlah data payed_buy_fee_sum pada tabel dws_shops guna 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;
  4. Lihat data tabel dws_users dan dws_shops.

    Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih Development > Scripts. Pada tab Scripts, klik + > New Script. 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;

    image

    -- Lihat data tabel dws_shops.
    SELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;

    截屏2024-09-02 14

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.

  1. Masukkan data berikut ke database order_dw 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');
  2. Lihat data tabel dws_users dan dws_shops.Pada menu navigasi kiri Konsol Pengembangan Realtime Compute for Apache Flink, pilih Development > Scripts. Pada tab Scripts, klik + > New Script. Pada editor SQL, salin dan tempel kode berikut, pilih kodenya, lalu Run:

    • Tabel dws_users

      SELECT * FROM paimoncatalog.order_dw.dws_users ORDER BY user_id;

      截屏2024-09-02 15

    • Tabel dws_shops

      SELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;

      截屏2024-09-02 15

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:

  • bucket: nama bucket OSS yang Anda buat.

  • object: path tempat data Anda disimpan.

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;

image

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;;

image

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;

image