Topik ini menjelaskan cara membuat tabel materialized, melakukan pengisian ulang data, mengubah kesegaran data, dan melihat alur data.
Batasan
-
Fitur ini hanya tersedia di Ververica Runtime (VVR) 8.0.10 dan versi yang lebih baru.
-
Tabel materialized hanya dapat dibuat dalam katalog Apache Paimon yang menggunakan Filesystem atau DLF untuk penyimpanan metadata. Katalog Apache Paimon kustom tidak didukung.
-
Anda harus memiliki izin untuk mengembangkan dan menyebarluaskan pekerjaan. Untuk informasi selengkapnya, lihat otorisasi konsol pengembangan.
-
Objek temporary, seperti tabel temporary, user-defined function temporary, dan view temporary, tidak didukung.
Buat materialized table
Sintaks
CREATE MATERIALIZED TABLE [catalog_name.][db_name.]table_name
-- Primary key constraint
[([CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED)]
[COMMENT table_comment]
-- Partition key
[PARTITIONED BY (partition_column_name1, partition_column_name2, ...)]
-- With options
[WITH (key1=val1, key2=val2, ...)]
-- Data freshness
FRESHNESS = INTERVAL '<num>' { SECOND | MINUTE | HOUR | DAY }
-- Refresh mode
[REFRESH_MODE = { CONTINUOUS | FULL }]
AS <select_statement>
Parameter
|
Parameter |
Wajib |
Deskripsi |
|
FRESHNESS |
Ya |
Kesegaran data dari tabel materialized, yang menentukan latensi maksimum yang diizinkan untuk pembaruan data dari tabel sumber. Catatan
|
|
AS <select_statement> |
Ya |
Menentukan kueri yang mengisi tabel materialized. Tabel hulu dapat berupa tabel materialized, tabel biasa, atau view. Pernyataan |
|
PRIMARY KEY |
Tidak |
Kolom opsional yang secara unik mengidentifikasi setiap baris dalam tabel. Kolom-kolom ini tidak boleh berisi nilai null. |
|
PARTITIONED BY |
Tidak |
Kolom opsional yang digunakan untuk mempartisi tabel materialized. |
|
WITH Options |
Tidak |
Menentukan properti tabel dan parameter format waktu untuk kolom partisi. Sebagai contoh, Anda dapat mengatur parameter format waktu untuk kolom partisi dengan |
|
REFRESH_MODE |
Tidak |
Menentukan mode refresh untuk tabel materialized. Mode refresh yang ditentukan memiliki prioritas lebih tinggi daripada mode yang secara otomatis disimpulkan oleh framework berdasarkan kesegaran data. Hal ini memungkinkan Anda menangani skenario tertentu.
|
Prosedur
-
Masuk ke Konsol Realtime Compute for Apache Flink.
-
Pada kolom Actions ruang kerja target, klik Console.
-
Di panel navigasi sebelah kiri, klik Catalogs, lalu klik katalog Apache Paimon target.
-
Klik database target, lalu klik Create Materialized Table.
Asumsikan Anda memiliki tabel sumber bernama
ordersdengan primary keyorder_id, nama kategoriorder_name, dan bidang tanggalds. Contoh berikut menunjukkan cara membuat tabel materialized berdasarkan tabel ini:-
Buat tabel materialized
mt_orderberdasarkan tabelorders. Kueri memilih semua kolom, dan kesegaran data diatur menjadi 5 detik.CREATE MATERIALIZED TABLE mt_order FRESHNESS = INTERVAL '5' SECOND AS SELECT * FROM `paimon`.`db`.`orders` ; -
Buat tabel materialized
mt_idberdasarkan tabel materializedmt_order. Kueri memilihorder_iddandssebagai kolom tabel, menetapkanorder_idsebagai primary key,dssebagai kolom partisi, dan kesegaran data menjadi 30 menit.CREATE MATERIALIZED TABLE mt_id ( PRIMARY KEY (order_id) NOT ENFORCED ) PARTITIONED BY(ds) FRESHNESS = INTERVAL '30' MINUTE AS SELECT order_id,ds FROM mt_order ; -
Buat tabel materialized
mt_dsberdasarkan tabel materializedmt_order, dan tentukandate-formatter(format waktu) untuk kolom partisids. Setiap kali penjadwalan dijalankan, waktu penjadwalan dikurangi kesegaran dikonversi menjadi nilai partisidsyang sesuai. Misalnya, jika kesegaran data diatur menjadi 1 jam dan waktu penjadwalan adalah2024-01-01 00:00:00, nilai ds yang dihitung adalah 20231231, dan hanya data dalam partisids = '20231231'yang direfresh. Jika waktu penjadwalan adalah2024-01-01 01:00:00, nilai ds yang dihitung adalah 20240101, dan data dalam partisids = '20240101'yang direfresh.CREATE MATERIALIZED TABLE mt_ds PARTITIONED BY(ds) WITH ( 'partition.fields.ds.date-formatter' = 'yyyyMMdd' ) FRESHNESS = INTERVAL '1' HOUR AS SELECT order_id,order_name,ds FROM mt_order ;Catatan-
Dalam
partition.fields.#.date-formatter, placeholder#harus merupakan kolom partisi valid bertipe STRING. -
Opsi
partition.fields.#.date-formattermenentukan format partisi waktu untuk tabel materialized. Placeholder#mewakili nama kolom partisi bertipe string. Informasi ini memungkinkan sistem mengidentifikasi partisi mana yang harus direfresh selama pembaruan terjadwal.
-
-
-
Mulai atau hentikan pembaruan tabel materialized.
-
Klik tabel materialized target di bawah katalognya.
-
Di pojok kanan atas, klik Start atau Stop.
CatatanJika Anda menghentikan pembaruan yang sedang berlangsung, pekerjaan akan menyelesaikan siklus pembaruan saat ini sebelum berhenti.
-
-
Lihat detail pekerjaan tabel materialized.
Pada tab Table Schema, di bagian Basic Information, klik ID pekerjaan di samping Latest Job atau Workflow untuk melihat detailnya.
Memodifikasi kueri tabel materialized
Batasan
-
Anda hanya dapat memodifikasi kueri untuk tabel materialized yang dibuat di VVR 11.1 atau versi yang lebih baru.
-
Saat memodifikasi kueri, Anda hanya dapat menambahkan kolom dan memodifikasi logika komputasi. Anda tidak dapat mengubah urutan kolom yang sudah ada atau memodifikasi definisinya.
Operasi
Didukung
Deskripsi
Tambahkan kolom baru
Ya
Anda dapat menambahkan kolom baru ke skema sambil mempertahankan urutan kolom yang sudah ada.
Modifikasi logika komputasi untuk kolom yang sudah ada (tanpa mengubah nama atau tipe datanya)
Ya
Anda dapat memodifikasi logika komputasi, tetapi nama kolom dan tipe data harus tetap sama.
Ubah urutan kolom yang sudah ada
Tidak
Urutan kolom bersifat tetap. Untuk mengubahnya, Anda harus menghapus dan membuat ulang tabel materialized.
Modifikasi nama atau tipe data kolom yang sudah ada
Tidak
Anda harus menghapus dan membuat ulang tabel materialized.
Contoh
-
Klik Edit Table dan modifikasi kueri. Kode berikut memberikan contohnya:
ALTER MATERIALIZED TABLE `paimon`.`default`.`mt-orders` AS SELECT *, price * quantity AS total_price FROM orders WHERE price * quantity > 1000 ; -
Klik Preview untuk melihat perbandingan sebelum dan sesudah.
Kotak dialog Materialized Table Change Details menampilkan perubahan berikut: di area SQL Statement, pernyataan SQL berubah dari
SELECT *yang asli menjadiSELECT *, <code data-tag="code" id="code_ee0dc0cd25b8">orders.price*orders.quantityAStotal_price, dan menyertakan kondisi WHERE baru<code data-tag="code" id="code_18ba1c5cea54">orders.price*orders.quantity> 1000. Di area Table Columns, bidangtotal_price(DOUBLE) yang baru disorot dengan warna hijau. -
Klik OK. Anda dapat melihat kolom yang baru ditambahkan dan logika kueri yang dimodifikasi pada tab Table Schema.
Penambahan kolom biasanya tidak memengaruhi konsumen hilir. Namun, jika pekerjaan hilir bergantung pada parsing dinamis (seperti SELECT * atau pemetaan bidang otomatis) saat mengonsumsi data dari tabel materialized hulu, pekerjaan tersebut mungkin gagal atau melaporkan error ketidaksesuaian format data. Kami menyarankan agar Anda menghindari parsing dinamis, menggunakan daftar kolom tetap, dan segera memperbarui skema tabel hilir setiap kali skema hulu berubah.
Pembaruan inkremental
Batasan
Fitur ini hanya tersedia di VVR 8.0.11 dan versi yang lebih baru.
Mode pembaruan
Tabel materialized mendukung tiga mode pembaruan: streaming, batch penuh, dan batch inkremental.
Mode ditentukan oleh pengaturan kesegaran data. Kesegaran kurang dari 30 menit mengaktifkan mode streaming, sedangkan 30 menit atau lebih mengaktifkan mode batch. Dalam mode batch, engine secara otomatis memutuskan antara pembaruan penuh atau inkremental. Pembaruan inkremental hanya menghitung data yang berubah sejak pembaruan terakhir dan menggabungkannya ke dalam tabel materialized. Pembaruan penuh menghitung data untuk seluruh tabel atau partisi dan menimpa data yang ada di tabel materialized. Dalam mode batch, engine memprioritaskan pembaruan inkremental dan hanya kembali ke pembaruan penuh jika pembaruan inkremental tidak memungkinkan.
Kondisi pembaruan inkremental
Pembaruan inkremental hanya dilakukan jika tabel materialized memenuhi semua kondisi berikut:
-
Parameter
partition.fields.#.date-formattertidak dikonfigurasi dalam definisi tabel. -
Tabel sumber tidak memiliki primary key yang ditentukan.
-
Kueri dalam definisi tabel materialized mendukung pembaruan inkremental sebagaimana dijelaskan dalam tabel berikut:
Klausa SQL
Dukungan
SELECT
Didukung untuk pemilihan kolom dan ekspresi fungsi skalar, termasuk user-defined function. Fungsi agregat tidak didukung.
FROM
Didukung untuk nama tabel dan subkueri.
WITH
Didukung untuk Common Table Expressions (CTEs).
WHERE
Didukung untuk kondisi filter yang mencakup berbagai ekspresi fungsi skalar, termasuk user-defined function. Subkueri seperti
WHERE [NOT] EXISTS <subquery>danWHERE <column> [NOT] IN <subquery>tidak didukung.UNION
Hanya
UNION ALLyang didukung.JOIN
-
INNER JOINdidukung. -
LEFT/RIGHT/FULL [OUTER] JOINtidak didukung, kecuali dalam kasus khususLATERAL JOINdan lookup join yang dijelaskan di bawah. -
[LEFT [OUTER]] JOIN LATERALdengan ekspresi fungsi tabel (termasuk user-defined function) didukung. -
Untuk lookup join, hanya
A [LEFT [OUTER]] JOIN B FOR SYSTEM_TIME AS OF PROCTIME()yang didukung.
Catatan-
Join implisit tanpa kata kunci
JOIN, sepertiSELECT * FROM a, b WHERE a.id = b.id, didukung. -
Komputasi inkremental untuk
INNER JOINtetap membaca seluruh data dari kedua tabel sumber.
GROUP BY
Tidak didukung.
-
Contoh pembaruan inkremental
Contoh 1: Proses data dari tabel sumber orders menggunakan fungsi skalar.
CREATE MATERIALIZED TABLE mt_shipped_orders (
PRIMARY KEY (order_id) NOT ENFORCED
)
FRESHNESS = INTERVAL '30' MINUTE
AS
SELECT
order_id,
COALESCE(customer_id, 'Unknown') AS customer_id,
CAST(order_amount AS DECIMAL(10, 2)) AS order_amount,
CASE
WHEN status = 'shipped' THEN 'Completed'
WHEN status = 'pending' THEN 'In Progress'
ELSE 'Unknown'
END AS order_status,
DATE_FORMAT(order_ts, 'yyyyMMdd') AS order_date,
UDSF_ProcessFunction(notes) AS notes
FROM
orders
WHERE
status = 'shipped';
Contoh 2: Perkaya data dari tabel sumber orders menggunakan LATERAL JOIN dan lookup join.
CREATE MATERIALIZED TABLE mt_enriched_orders (
PRIMARY KEY (order_id, order_tag) NOT ENFORCED
)
FRESHNESS = INTERVAL '30' MINUTE
AS
WITH o AS (
SELECT
order_id,
product_id,
quantity,
proc_time,
e.tag AS order_tag
FROM
orders,
LATERAL TABLE(UDTF_StringSplitFunction(tags, ',')) AS e(tag))
SELECT
o.order_id,
o.product_id,
p.product_name,
p.category,
o.quantity,
p.price,
o.quantity * p.price AS total_amount,
order_tag
FROM o
LEFT JOIN
product_info FOR SYSTEM_TIME AS OF PROCTIME() AS p
ON
o.product_id = p.product_id;
Pengisian ulang data
Sebelumnya, memperbaiki hasil pemrosesan aliran untuk data historis memerlukan pengembangan pekerjaan batch terpisah. Dengan menggunakan tabel materialized, Anda dapat langsung mengisi ulang partisi data historis. Pendekatan ini menyatukan pemrosesan batch dan streaming, sehingga mengurangi biaya pengembangan dan O&M.
-
Klik tabel materialized target di bawah katalognya.
-
Pada tab Data Information, isi ulang data.
Jika Anda menentukan kolom partisi saat membuat tabel materialized, tabel tersebut adalah tabel partisi. Jika tidak, tabel tersebut adalah tabel non-partisi.
Tabel partisi
Di bagian Partitions, klik Trigger Update jika ini pertama kalinya Anda mengisi ulang data atau jika partisi yang diperlukan belum ada. Jika partisi sudah ada, Anda dapat memilih partisi tertentu dan klik Refresh di kolom Actions.
Klik tab Data Information dan lakukan operasi di area Data Partitions di bagian bawah halaman.
Parameter
-
Kolom partisi: Kolom partisi tabel. Misalnya, jika Anda memasukkan
20241201, semua data dalam partisids=20241201akan diisi ulang. -
Nama Tugas: Nama tugas pengisian ulang data.
-
Rentang Pembaruan (Opsional): Menentukan apakah akan melakukan pembaruan kaskade ke tabel materialized hilir. Dimulai dari tabel saat ini, semua tabel materialized dalam alur data akan diperbarui. Kedalaman maksimum alur data hilir yang didukung adalah 6 level.
Catatan-
Saat memperbarui tabel partisi, tabel materialized hilir harus memiliki kolom partisi yang persis sama dengan tabel awal. Jika tidak, operasi pembaruan akan gagal.
-
Jika pembaruan gagal untuk tabel materialized apa pun dalam alur data, semua node hilir berikutnya juga akan gagal.
-
-
Target Penyebaran: Anda dapat memilih antrian atau kluster session. Pilihan default adalah
default-queue.
Tabel non-partisi
Di bagian Data Detail, klik Refresh.
Parameter
-
Nama Tugas: Nama tugas pengisian ulang data.
-
Rentang Pembaruan: Opsi ini tidak tersedia untuk tabel non-partisi.
Catatan-
Selama pembaruan, data hilir direfresh sepenuhnya.
-
Jika pembaruan gagal untuk tabel materialized apa pun dalam alur data, semua node hilir berikutnya juga akan gagal.
-
Pembaruan kaskade tidak didukung jika tabel awal adalah tabel non-partisi yang diperbarui oleh pekerjaan streaming.
-
-
Target Penyebaran: Anda dapat memilih antrian atau kluster session. Pilihan default adalah
default-queue.
-
-
Pengisian ulang terjadwal dan batch.
Anda dapat menggunakan Task Orchestration untuk membuat alur kerja bagi tabel materialized agar menjalankan pekerjaan pengisian ulang secara terjadwal. Anda juga dapat menggunakan fitur pengisian ulang data alur kerja untuk mengisi ulang data secara massal dalam rentang waktu tertentu.
Ubah kesegaran data
-
Di katalog yang sesuai, klik database materialized table, lalu pilih materialized table target.
-
Di pojok kanan atas, klik Modify Freshness.
-
Jika tabel materialized tidak memiliki primary key, Anda tidak dapat mengganti metode pembaruannya antara streaming dan batch. Sistem menggunakan pekerjaan streaming untuk nilai kesegaran data di bawah 30 menit dan pekerjaan batch untuk nilai 30 menit atau lebih. Oleh karena itu, mengubah kesegaran data melewati ambang batas 30 menit ini tidak diizinkan untuk tabel tanpa primary key.
-
Jika tabel hulu adalah tabel materialized, pastikan kesegaran data tabel hilir merupakan kelipatan bilangan bulat positif dari kesegaran data tabel hulu.
-
Kesegaran data maksimum adalah satu hari.
-
Lihat alur data
Di panel navigasi sebelah kiri, pilih untuk membuka halaman alur data untuk tabel materialized. Di halaman ini, Anda dapat melihat hubungan alur data antara semua tabel materialized. Anda juga dapat melakukan operasi seperti Start/Stop Update dan Modify Freshness langsung pada tabel materialized. Klik Details untuk membuka halaman detail tabel materialized yang sesuai.
Klik node tabel materialized untuk memperluas panel detailnya. Panel tersebut menampilkan nilai Data Freshness, Latest Update Time, dan Materialized Table Status, serta menyediakan aksi Trigger Update.
Dokumen terkait
-
Untuk pengenalan tabel materialized, lihat Kelola tabel materialized.
-
Untuk mempelajari cara membangun pipeline analitis di danau data terpadu dengan pemrosesan batch dan streaming terpadu berdasarkan Paimon dan tabel materialized, serta cara beralih dari mode batch ke mode streaming dengan memodifikasi kesegaran data untuk pembaruan real-time, lihat Panduan cepat: Bangun danau data terpadu dengan pemrosesan batch dan streaming terpadu menggunakan tabel materialized.