Topik ini menjelaskan cara menggunakan Paimon connector untuk streaming data lakehouse. Untuk hasil terbaik, kami merekomendasikan penggunaan connector ini bersama Paimon Catalog.
Informasi latar belakang
Apache Paimon adalah format penyimpanan data lake terpadu untuk aliran (streaming) dan batch yang mendukung penulisan throughput tinggi serta kueri latensi rendah. Paimon terintegrasi dengan baik dengan mesin komputasi populer yang tersedia di Alibaba Cloud E-MapReduce, seperti Flink, Spark, Hive, dan Trino. Anda dapat menggunakan Apache Paimon untuk membangun data lake secara cepat di HDFS atau OSS dan menghubungkannya ke mesin komputasi tersebut guna melakukan analitik data lake. Untuk informasi lebih lanjut, lihat Apache Paimon.
|
Kategori |
Deskripsi |
|
Jenis yang didukung |
Tabel sumber, tabel dimensi, tabel hasil, dan target untuk ingesti data |
|
Mode eksekusi |
Mode streaming dan mode batch |
|
Format data |
Tidak didukung |
|
Metrik pemantauan |
Tidak ada |
|
Jenis API |
SQL dan YAML untuk ingesti data |
|
Pembaruan dan penghapusan pada tabel hasil |
Ya |
Fitur utama
Apache Paimon menyediakan kemampuan inti berikut:
-
Membangun data lake ringan dan berbiaya rendah di HDFS atau object storage.
-
Membaca dan menulis dataset skala besar dalam mode streaming maupun batch.
-
Menjalankan kueri batch dan OLAP dengan freshness data dari hitungan menit hingga detik.
-
Mengingesti dan menghasilkan data inkremental, berfungsi sebagai lapisan penyimpanan baik untuk gudang data offline tradisional maupun gudang data streaming modern.
-
Melakukan pre-agregasi data untuk mengurangi biaya penyimpanan dan beban komputasi downstream.
-
Mengakses versi historis data.
-
Menyaring data secara efisien.
-
Mendukung evolusi skema.
Batasan dan rekomendasi
-
Paimon connector memerlukan mesin komputasi Flink VVR 6.0.6 atau versi yang lebih baru.
-
Tabel berikut mencantumkan kompatibilitas versi antara Paimon dan VVR.
Versi Apache Paimon
VVR
1.3.1
11.5, 11.6, 11.7, 11.8
1.3
11.4
1.2
11.2, 11.3
1.1
11.1
1.0
8.0.11
-
Rekomendasi penyimpanan untuk penulisan konkuren
Saat beberapa job melakukan penulisan konkuren ke tabel Paimon yang sama, penggunaan penyimpanan OSS standar (oss://) kadang-kadang dapat menyebabkan konflik commit atau kegagalan job karena keterbatasan operasi file atomik.
Untuk penulisan yang stabil dan konsisten, gunakan layanan metadata atau penyimpanan yang menyediakan jaminan atomic kuat. Opsi yang disarankan adalah Data Lake Formation (DLF), yang menawarkan manajemen terpadu untuk metadata dan penyimpanan Paimon. Sebagai alternatif, Anda dapat menggunakan OSS-HDFS atau HDFS.
-
Cara perubahan konfigurasi diterapkan
Perubahan pada parameter konfigurasi tabel Paimon hanya berlaku setelah Anda me-restart job terkait. Job yang sedang berjalan tidak memuat perubahan ini secara dinamis.
-
Reklamasi fisik tertunda untuk partisi yang dihapus
Saat Anda menjalankan operasi DROP PARTITION, sistem tidak langsung menghapus file data fisik yang mendasarinya.
Operasi ini melakukan penghapusan logis. Paimon hanya menghapus metadata partisi target dari snapshot terbaru. Karena Paimon mendukung fitur time travel, snapshot historis masih mereferensikan file data partisi tersebut. File data fisik akan dihapus secara permanen hanya setelah semua snapshot historis yang mereferensikan partisi tersebut mencapai batas retensi mereka dan dibersihkan oleh mekanisme kedaluwarsa snapshot.
SQL
Gunakan Paimon connector dalam job SQL sebagai tabel sumber atau tabel sink.
Sintaks
-
Jika Anda membuat tabel Paimon di Paimon Catalog, Anda tidak perlu menentukan parameter
connector. Sintaksnya sebagai berikut:CREATE TABLE `<YOUR-PAIMON-CATALOG>`.`<YOUR-DB>`.paimon_table ( id BIGINT, data STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( ... );CatatanJika Anda telah membuat tabel Paimon di Paimon Catalog, Anda dapat langsung menggunakannya.
-
Jika Anda membuat tabel temporary Paimon di Catalog lain, Anda harus menentukan parameter 'connector' dan 'path'. Sintaksnya sebagai berikut:
CREATE TEMPORARY TABLE paimon_table ( id BIGINT, data STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'paimon', 'path' = '<path-to-paimon-table-files>', 'auto-create' = 'true', -- Jika file data tabel Paimon tidak ada di path yang ditentukan, file tersebut akan dibuat secara otomatis. ... );Catatan-
Contoh path:
'path' = 'oss://<bucket>/test/order.db/orders'. Jangan menghilangkan akhiran.db. Paimon mengandalkan akhiran ini untuk mengidentifikasi database. -
Beberapa job yang menulis ke tabel yang sama harus menggunakan konfigurasi path yang sama.
-
Jika dua konfigurasi path berbeda, Paimon tidak menganggapnya sebagai tabel yang sama. Meskipun path fisiknya sama, konfigurasi Catalog yang tidak konsisten dapat menyebabkan konflik penulisan konkuren, kegagalan operasi compact, dan kehilangan data. Misalnya, Paimon menganggap
oss://b/testdanoss://b/test/sebagai tabel berbeda karena adanya garis miring di akhir, meskipun keduanya mungkin mengarah ke lokasi fisik yang sama.
-
Parameter WITH
|
Parameter |
Deskripsi |
Tipe |
Wajib |
Bawaan |
Keterangan |
|
connector |
Menentukan connector untuk tabel. |
String |
Tidak |
Tidak ada |
|
|
path |
Path penyimpanan tabel. |
String |
Tidak |
Tidak ada |
|
|
auto-create |
Menentukan apakah akan membuat file tabel secara otomatis jika file tersebut tidak ada di path yang ditentukan. |
Boolean |
Tidak |
false |
Nilai yang valid:
|
|
file.format |
Format file data. |
String |
Tidak |
parquet |
Nilai yang valid:
|
|
bucket |
Jumlah bucket per partisi. |
Integer |
Tidak |
1 |
Paimon mendistribusikan data ke bucket berdasarkan Catatan
Kami merekomendasikan agar setiap bucket berisi kurang dari 5 GB data. |
|
bucket-key |
Kolom yang digunakan sebagai kunci bucket. |
String |
Tidak |
Tidak ada |
Menentukan kolom yang digunakan untuk mendistribusikan data ke bucket. Pisahkan beberapa nama kolom dengan koma (,). Contohnya, Catatan
|
|
changelog-producer |
Mekanisme produksi changelog. |
String |
Tidak |
none |
Paimon dapat menghasilkan changelog lengkap untuk aliran input apa pun, artinya setiap catatan
Untuk informasi lebih lanjut tentang cara memilih produsen changelog, lihat Changelog Production. |
|
full-compaction.delta-commits |
Jumlah maksimum commit antara dua full compaction berturut-turut. |
Integer |
Tidak |
Tidak ada |
Menentukan jumlah maksimum commit snapshot yang diizinkan sebelum memicu full compaction. |
|
lookup.cache-max-memory-size |
Ukuran cache memori untuk tabel dimensi Paimon. |
String |
Tidak |
256 MB |
Parameter ini mengontrol ukuran cache baik untuk lookup tabel dimensi maupun produsen changelog |
|
merge-engine |
Mekanisme penggabungan catatan dengan kunci primer yang sama. |
String |
Tidak |
deduplicate |
Nilai yang valid:
Untuk analisis mendetail tentang mesin penggabungan, lihat Merge Engine. |
|
partial-update.ignore-delete |
Menentukan apakah akan mengabaikan pesan penghapusan (-D). |
Boolean |
Tidak |
false |
Nilai yang valid:
Catatan
|
|
ignore-delete |
Menentukan apakah akan mengabaikan pesan penghapusan (-D). |
Boolean |
Tidak |
false |
Nilai yang valid sama dengan nilai untuk partial-update.ignore-delete. Catatan
|
|
partition.default-name |
Nama partisi bawaan. |
String |
Tidak |
__DEFAULT_PARTITION__ |
Nama partisi yang digunakan ketika nilai kolom partisi adalah null atau string kosong. |
|
partition.expiration-check-interval |
Frekuensi sistem memeriksa partisi yang kedaluwarsa. |
String |
Tidak |
1h |
Untuk detailnya, lihat Cara mengonfigurasi kedaluwarsa partisi otomatis. |
|
partition.expiration-time |
Durasi setelah partisi kedaluwarsa. |
String |
Tidak |
Tidak ada |
Partisi kedaluwarsa ketika usianya melebihi nilai ini. Secara bawaan, partisi tidak pernah kedaluwarsa. Sistem menghitung usia partisi dari nilai partisinya. Untuk detailnya, lihat Cara mengonfigurasi kedaluwarsa partisi otomatis. |
|
partition.timestamp-formatter |
String format untuk mengonversi string waktu menjadi timestamp. |
String |
Tidak |
Tidak ada |
Menentukan format untuk mengekstraksi usia partisi dari nilai partisi. Untuk detailnya, lihat Cara mengonfigurasi kedaluwarsa partisi otomatis. |
|
partition.timestamp-pattern |
String format untuk mengonversi nilai partisi menjadi string waktu. |
String |
Tidak |
Tidak ada |
Menentukan pola untuk mengekstraksi string waktu dari nilai partisi. Untuk detailnya, lihat Cara mengonfigurasi kedaluwarsa partisi otomatis. |
|
scan.bounded.watermark |
Nilai watermark yang menandai akhir pemindaian. Tabel sumber berhenti menghasilkan data ketika watermark-nya melebihi nilai ini. |
Long |
Tidak |
Tidak ada |
Tidak berlaku |
|
scan.mode |
Menentukan posisi konsumsi untuk tabel sumber Paimon. |
String |
Tidak |
bawaan |
Untuk detailnya, lihat Cara menetapkan posisi konsumsi untuk tabel sumber Paimon. |
|
scan.snapshot-id |
Menentukan snapshot tempat tabel sumber Paimon mulai mengonsumsi. |
Integer |
Tidak |
Tidak ada |
Untuk detailnya, lihat Cara menetapkan posisi konsumsi untuk tabel sumber Paimon. |
|
scan.timestamp-millis |
Menentukan titik waktu tempat tabel sumber Paimon mulai mengonsumsi. |
Integer |
Tidak |
Tidak ada |
Untuk detailnya, lihat Cara menetapkan posisi konsumsi untuk tabel sumber Paimon. |
|
snapshot.num-retained.max |
Jumlah maksimum snapshot terbaru yang dipertahankan. |
Integer |
Tidak |
2147483647 |
Kedaluwarsa snapshot dipicu jika salah satu kondisi ini atau kondisi |
|
snapshot.num-retained.min |
Jumlah minimum snapshot terbaru yang dipertahankan. |
Integer |
Tidak |
10 |
Tidak berlaku |
|
snapshot.time-retained |
Periode retensi untuk snapshot. |
String |
Tidak |
1h |
Kedaluwarsa snapshot dipicu jika salah satu kondisi ini atau kondisi |
|
write-mode |
Mode penulisan untuk tabel Paimon. |
String |
Tidak |
change-log |
Nilai yang valid:
Untuk informasi lebih lanjut tentang mode penulisan, lihat Write Mode. |
|
scan.infer-parallelism |
Menentukan apakah akan secara otomatis menginfer paralelisme untuk tabel sumber Paimon. |
Boolean |
Tidak |
true |
Nilai yang valid:
|
|
scan.parallelism |
Paralelisme untuk tabel sumber Paimon. |
Integer |
Tidak |
Tidak ada |
Catatan
Parameter ini diabaikan jika Mode diatur ke Mode di tab job. |
|
sink.parallelism |
Paralelisme untuk tabel sink Paimon. |
Integer |
Tidak |
Tidak ada |
Catatan
Parameter ini diabaikan jika Mode diatur ke Mode di tab job. |
|
sink.clustering.by-columns |
Menentukan kolom clustering untuk penulisan ke tabel sink Paimon. |
String |
Tidak |
Tidak ada |
Untuk tabel append-only Paimon (tabel tanpa kunci primer), parameter ini mengaktifkan penulisan terklaster dalam job batch. Proses ini meningkatkan performa kueri dengan mengelompokkan data pada kolom yang ditentukan. Pisahkan beberapa nama kolom dengan koma (,), contohnya, Untuk informasi lebih lanjut tentang clustering, lihat Dokumentasi Resmi Apache Paimon. |
|
sink.delete-strategy |
Menentukan strategi validasi untuk memastikan sistem menangani pesan retraction (-D/-U) dengan benar. |
Enum |
Tidak |
NONE |
Nilai yang valid dan perilaku yang diharapkan dari operator sink saat menangani pesan retraction:
Catatan
|
|
blob-as-descriptor |
Menentukan apakah akan mengeluarkan byte descriptor Blob saat kolom Blob dibaca. |
Boolean |
Tidak |
false |
Saat parameter ini diatur ke true, kueri mengembalikan byte serialisasi BlobDescriptor alih-alih konten Blob aktual. Gunakan parameter ini bersama fungsi URL yang ditandatangani Blob. Anda tidak perlu mengatur parameter ini saat pembuatan tabel. Anda dapat mengaturnya secara dinamis saat membaca dengan menggunakan petunjuk SQL. Parameter ini didukung di VVR 11.9-preview1 dan versi yang lebih baru. Catatan
Hanya didukung di VVR 11.9-preview1 dan versi yang lebih baru. |
|
sink.existing-table.schema-check.enabled |
Menentukan apakah akan memeriksa kompatibilitas kunci primer, jumlah bidang, dan tipe bidang dari skema tabel saat menulis data ke tabel yang telah dibuat sebelumnya. |
Boolean |
Tidak |
true |
Hanya didukung di VVR 11.8 dan versi yang lebih baru. |
Untuk informasi lebih lanjut tentang opsi konfigurasi, lihat Dokumentasi Resmi Apache Paimon.
Parameter tabel vektor
Parameter berikut hanya didukung di VVR 11.8 dan versi yang lebih baru.
|
Parameter |
Deskripsi |
Tipe data |
Wajib |
Nilai bawaan |
Keterangan |
|
|
Jumlah kluster IVF yang diproses selama pencarian. |
Integer |
Tidak |
16 |
Nilai yang lebih tinggi umumnya meningkatkan recall tetapi meningkatkan latensi. |
|
|
Mengambil kandidat IVF top_k × refine_factor dan melakukan re-ranking menggunakan vektor asli yang disimpan di tabel Paimon. |
Float |
Tidak |
Nonaktif |
Nonaktif secara bawaan untuk semua varian IVF. Paling cocok untuk indeks terkompresi (seperti ivf-pq dan ivf-hnsw-sq) dalam skenario di mana recall diprioritaskan daripada latensi. |
|
|
Lebar pencarian HNSW selama pencarian. |
Integer |
Tidak |
0 |
Nilai yang lebih tinggi umumnya meningkatkan recall tetapi meningkatkan latensi. 0 menunjukkan bahwa bawaan pustaka native digunakan. |
|
|
Ukuran daftar pencarian Lumina DiskANN. |
Integer |
Tidak |
max(1.5 × top_k, 16) |
Nilai yang lebih tinggi umumnya meningkatkan recall tetapi meningkatkan latensi. |
|
|
Lebar beam pencarian Lumina DiskANN. |
Integer |
Tidak |
4 |
— |
|
|
Jumlah pencarian Lumina paralel. |
Integer |
Tidak |
5 |
— |
Detail fitur
Freshness dan konsistensi data
Tabel sink Paimon menggunakan protokol two-phase commit untuk meng-commit data selama setiap checkpoint job Flink. Oleh karena itu, freshness data ditentukan oleh interval checkpoint job Flink. Setiap commit menghasilkan hingga dua snapshot.
Saat dua job Flink menulis ke tabel Paimon yang sama secara konkuren, jika job tersebut menulis ke bucket yang berbeda, mereka mencapai konsistensi serializable. Jika job tersebut menulis ke bucket yang sama, mereka hanya mencapai isolasi snapshot. Artinya, data tabel mungkin merupakan campuran hasil dari kedua job tersebut, tetapi tidak terjadi kehilangan data.
Mesin penggabungan
Saat tabel sink Paimon menerima beberapa catatan dengan kunci primer yang sama, tabel tersebut menggabungkannya menjadi satu catatan untuk menjaga keunikan. Anda dapat mengontrol perilaku ini dengan mengatur parameter merge-engine. Tabel berikut menjelaskan mesin penggabungan yang tersedia.
|
Merge engine |
Deskripsi |
|
Deduplicate |
Mesin deduplikasi adalah bawaan. Untuk beberapa catatan dengan kunci primer yang sama, tabel sink Paimon hanya menyimpan catatan terbaru dan membuang yang lainnya. Catatan
Jika catatan terbaru adalah pesan penghapusan, semua catatan dengan kunci primer tersebut dibuang. |
|
Partial Update |
Mesin partial update memungkinkan Anda membangun catatan lengkap dengan memperbarui secara inkremental menggunakan beberapa pesan. Saat catatan baru dengan kunci primer yang sama tiba, nilai non-null-nya menimpa bidang yang sesuai dalam catatan yang ada. Mesin mengabaikan bidang yang bernilai null dalam catatan baru dan mempertahankan nilai yang ada. Sebagai contoh, asumsikan tabel sink Paimon menerima tiga catatan berikut secara berurutan:
Jika kolom pertama adalah kunci primer, catatan gabungan akhirnya adalah <1, 25.2, 10, 'This is a book'>. Catatan
|
|
Aggregation |
Dalam beberapa kasus penggunaan, Anda mungkin hanya memerlukan nilai agregat catatan. Mesin agregasi menggabungkan catatan yang memiliki kunci primer yang sama menggunakan fungsi agregat yang Anda tentukan. Untuk setiap kolom non-kunci-primer, Anda harus menentukan fungsi agregat menggunakan opsi
Kolom
Catatan
|
Produsen log perubahan
Atur parameter changelog-producer untuk mengonfigurasi Paimon agar menghasilkan changelog lengkap (di mana setiap catatan update_after memiliki catatan update_before yang sesuai) untuk aliran input apa pun. Tabel berikut menjelaskan produsen changelog yang tersedia. Untuk detail lebih lanjut, lihat dokumentasi resmi Apache Paimon.
|
Produsen |
Deskripsi |
|
None |
Saat Anda mengatur Sebagai contoh, jika konsumen downstream perlu menghitung jumlah suatu kolom dan hanya melihat nilai terbaru 5, konsumen tersebut tidak dapat menentukan cara memperbarui total. Jika nilai sebelumnya adalah 4, jumlah harus bertambah 1; jika nilai sebelumnya adalah 6, jumlah harus berkurang 1. Konsumen yang sensitif terhadap catatan Catatan
Jika konsumen downstream Anda, seperti database, tidak sensitif terhadap data |
|
Input |
Saat Anda mengatur Gunakan produsen ini hanya ketika aliran input itu sendiri sudah merupakan changelog lengkap, seperti data dari Change Data Capture (CDC). |
|
Lookup |
Saat Anda mengatur Dibandingkan dengan produsen Gunakan opsi ini untuk kasus penggunaan yang memerlukan freshness data tinggi (misalnya, tingkat menit). |
|
Full Compaction |
Saat Anda mengatur Dibandingkan dengan produsen Gunakan opsi ini untuk kasus penggunaan dengan persyaratan freshness data rendah (misalnya, tingkat jam). |
Mode penulisan
Tabel Paimon mendukung mode penulisan berikut.
|
Mode |
Deskripsi |
|
Change-log |
Mode penulisan |
|
Append-only |
Mode penulisan Untuk deskripsi mendetail tentang mode penulisan
|
Target untuk CTAS dan CDAS
Tabel Paimon mendukung sinkronisasi data real-time untuk tabel tunggal atau seluruh database. Perubahan skema di tabel upstream juga disinkronkan ke tabel Paimon secara real-time. Untuk detailnya, lihat Kelola tabel Paimon dan Kelola Paimon Catalog.
Pemangkasan baca Variant
Hanya membaca bidang Variant yang direferensikan oleh kueri, mengurangi overhead I/O dan memori.
Cara mengaktifkan
|
Parameter |
Nilai bawaan |
Deskripsi |
|
|
false |
Atur ke |
Prasyarat
-
Data Variant hanya dapat ditulis oleh VVR 11.6 atau versi yang lebih baru. Data Variant yang ditulis oleh versi sebelumnya atau mesin lain tidak mendukung pemangkasan baca.
-
Agar data upstream berlaku, Anda harus mengaktifkan opsi
variant.inferShreddingSchemauntuk mengaktifkan inferensi shredding Variant otomatis, atau secara eksplisit mengonfigurasivariant.shreddingSchemadan opsi terkait untuk menentukan anggota Variant mana yang cocok untuk shredding. Untuk informasi lebih lanjut, lihat Coreoptions.
Skenario yang didukung
Pemangkasan baca berlaku saat SQL mengakses bidang Variant dengan kunci string, misalnya:
SELECT v['a'] FROM t; -- Satu level
SELECT v['a']['b'] FROM t; -- Bersarang
SELECT v['a'], v['b'] FROM t; -- Beberapa bidang
SELECT id, v['a'] FROM t; -- Dicampur dengan kolom biasa
SELECT v['a'] + 1 FROM t; -- Direferensikan dalam ekspresi
Skenario yang tidak didukung
-
Mereferensikan seluruh Variant secara langsung, seperti
SELECT v FROM t. -
Menggabungkan referensi langsung dengan akses bidang, seperti
SELECT v, v['a'] FROM t. -
Akses indeks array, seperti
v[0]. -
Skema tabel berisi bidang Row bersarang. Dalam kasus ini, pemangkasan baca tidak berlaku untuk bidang Variant apa pun di tabel tersebut.
URL Blob yang Telah Ditandatangani
Paimon dapat menghasilkan URL HTTPS GET publik berumur pendek yang ditandatangani untuk data kolom Blob yang disimpan di OSS. Layanan eksternal, seperti layanan model multimodal, dapat menggunakan URL tersebut untuk mengambil konten file secara langsung, tanpa mentransfer byte file di dalam job Flink.
Batasan
-
Hanya mesin komputasi Flink VVR 11.9-preview1 dan versi yang lebih baru yang mendukung fitur ini.
-
Hanya tabel Paimon yang disimpan di OSS yang didukung. Jenis metastore tidak menjadi masalah. Jenis metastore Filesystem, DLF, dan lainnya semuanya didukung.
-
URL yang dihasilkan adalah URL HTTPS publik dan harus dikonsumsi oleh layanan eksternal yang memiliki akses jaringan publik. Jika katalog menggunakan titik akhir internal, kami merekomendasikan Anda mengubahnya ke titik akhir HTTPS publik agar layanan eksternal dapat mengakses URL yang dihasilkan.
-
Untuk membaca kolom Blob biasa, Anda harus terlebih dahulu mendapatkan byte BlobDescriptor dengan menggunakan opsi
blob-as-descriptor, lalu meneruskan byte tersebut ke fungsi.
Sintaks
descriptor_to_presigned_url(source_table, descriptor, validity)
try_descriptor_to_presigned_url(source_table, descriptor, validity)
Parameter
|
Parameter |
Tipe |
Deskripsi |
|
|
LITERAL STRING |
Tabel Paimon yang berisi kolom Blob, dalam format |
|
|
BYTES |
Byte BlobDescriptor, yang merupakan hasil membaca kolom Blob setelah |
|
|
INTERVAL |
Periode validitas URL. Nilainya harus berupa jumlah positif dalam detik, misalnya, |
Nilai kembalian
|
Type |
Deskripsi |
|
STRING |
URL HTTPS GET publik berumur pendek yang ditandatangani. |
Contoh
-- Gunakan petunjuk SQL untuk membaca kolom Blob sebagai descriptor dan menghasilkan URL yang ditandatangani
SELECT sys.descriptor_to_presigned_url(
'default.image_table',
image,
INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;
-- Mode toleransi kesalahan: mengembalikan NULL pada error tingkat baris, dengan perilaku lain tidak berubah
SELECT sys.try_descriptor_to_presigned_url(
'default.image_table',
image,
INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;
-
descriptor_to_presigned_urlmelemparkan exception pada error tingkat baris.try_descriptor_to_presigned_urlmengembalikan NULL pada error tingkat baris. -
URL yang dihasilkan tidak berisi ekstensi file, dan byte yang ditunjuk oleh URL identik dengan data Blob. Layanan downstream harus mengidentifikasi format file berdasarkan byte yang dikembalikan, bukan akhiran URL.
-
URL yang ditandatangani setara dengan kredensial akses sementara. Setelah URL yang ditandatangani dihasilkan, segera kirimkan ke layanan downstream. Jangan menuliskannya ke log atau tabel sink, dan jangan menyimpannya secara persisten.
Gunakan Paimon sebagai tabel dimensi
Tabel Paimon dapat digunakan sebagai tabel dimensi. Untuk sintaks JOIN, lihat Pernyataan JOIN tabel dimensi.
Secara bawaan, lookup memuat semua data di setiap instans paralel. Pendekatan ini hanya cocok untuk tabel dimensi kecil. Untuk tabel dimensi besar, gunakan solusi Shuffle Lookup yang dijelaskan di bawah.
Tabel dimensi terpartisi
Jika tabel dimensi Anda terpartisi dan Anda hanya memerlukan data dari satu atau dua partisi terbaru, Anda dapat menggunakan fitur pemuatan partisi dinamis:
SELECT * FROM T
JOIN DIM /*+ OPTIONS('lookup.dynamic-partition'='max_pt()', 'lookup.dynamic-partition.refresh-interval'='1 h') */
FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col = D.col;
|
Parameter |
Tipe data |
Nilai bawaan |
Deskripsi |
|
lookup.dynamic-partition |
String |
Tidak Berlaku |
|
|
lookup.dynamic-partition.refresh-interval |
Durasi |
1 h |
Interval sistem memeriksa pembaruan partisi di tabel dimensi. |
Tabel dimensi besar: tabel fixed-bucket
Hanya didukung di VVR 8.0.8 dan versi yang lebih baru. Untuk tabel fixed-bucket (bucket > 0), Anda dapat menggunakan Shuffle Lookup untuk mendistribusikan data berdasarkan kunci bucket di instans paralel, sehingga setiap instans hanya memuat data di bucket yang ditugaskan:
SELECT /*+ LOOKUP('table'='D', 'shuffle'='true') */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
-
Kunci join harus merupakan kunci bucket. Kunci bucket secara bawaan adalah kunci primer.
-
Hanya tabel fixed-bucket (bucket > 0) yang mendukung fitur ini.
Tabel dimensi besar: tabel non-fixed-bucket
Hanya didukung di VVR 8.0.10 dan versi yang lebih baru. Untuk tabel dynamic-bucket atau tabel append, Anda dapat menggunakan SHUFFLE_HASH atau REPLICATED_SHUFFLE_HASH sehingga setiap instans paralel membaca semua data tetapi hanya menyimpan bagian yang dibutuhkan:
-- Shuffle Hash
SELECT /*+ SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
-- Replicated Shuffle Hash
SELECT /*+ REPLICATED_SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
Untuk informasi lebih lanjut tentang SHUFFLE_HASH dan REPLICATED_SHUFFLE_HASH, lihat Pernyataan JOIN tabel dimensi.
Ingesti Data
Anda dapat menggunakan Paimon connector sebagai sink dalam job ingesti data YAML.
Sintaks
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: filesystem
catalog.properties.warehouse: /path/warehouse
Parameter
|
Parameter |
Deskripsi |
Wajib |
Tipe |
Bawaan |
Catatan |
|
type |
Jenis connector. |
Ya |
STRING |
Tidak ada |
Nilainya harus |
|
name |
Nama sink. |
Tidak |
STRING |
Tidak ada |
|
|
catalog.properties.metastore |
Jenis katalog Paimon. |
Tidak |
STRING |
filesystem |
Nilai yang valid:
|
|
catalog.properties.* |
Parameter untuk membuat katalog Paimon. |
Tidak |
STRING |
Tidak ada |
Untuk informasi lebih lanjut, lihat Kelola Paimon Catalog. |
|
table.properties.* |
Parameter untuk membuat tabel Paimon. |
Tidak |
STRING |
Tidak ada |
Untuk informasi lebih lanjut, lihat Opsi tabel Paimon. |
|
catalog.properties.warehouse |
Direktori root untuk penyimpanan file. |
Tidak |
STRING |
Tidak ada |
Parameter ini hanya berlaku saat |
|
commit.user-prefix |
Awalan username untuk meng-commit file data. |
Tidak |
STRING |
Tidak ada |
Catatan
Kami merekomendasikan menetapkan username berbeda untuk job berbeda. Hal ini mempermudah identifikasi job yang menyebabkan konflik commit. |
|
partition.key |
Kunci partisi untuk tabel terpartisi. |
Tidak |
STRING |
Tidak ada |
Tabel berbeda dipisahkan oleh |
|
sink.cross-partition-upsert.tables |
Daftar tabel yang memerlukan upsert lintas partisi, di mana kunci primer tidak mencakup semua kunci partisi. |
Tidak |
STRING |
Tidak ada |
Berlaku untuk tabel dengan pembaruan lintas partisi.
Penting
|
|
sink.commit.parallelism |
Menentukan paralelisme operator Commit. |
Tidak |
INTEGER |
Tidak ada |
Jika operator Commit menjadi bottleneck, gunakan parameter ini untuk meningkatkan paralelismenya dan meningkatkan performa. Parameter ini hanya didukung di Realtime Compute for Apache Flink 11.6 dan versi yang lebih baru. Catatan
Menyetel parameter ini mengubah paralelisme operator. Saat me-restart job stateful, Anda harus menentukan |
Gunakan katalog yang sudah ada
Mulai Realtime Compute for Apache Flink 11.5, Anda dapat langsung mereferensikan katalog Paimon bawaan dari halaman Data Management dalam job ingesti data Flink CDC. Hal ini mengurangi konfigurasi manual.
sink:
type: paimon
using.built-in-catalog: paimon_dlf_catalog
catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com
Job ingesti data dapat secara otomatis menggunakan kembali semua parameter dalam katalog Paimon. Hal ini setara dengan mengonfigurasi secara manual parameter yang diawali dengan catalog.properties. dalam job YAML.
Untuk mengganti parameter yang digunakan kembali secara otomatis, atur secara eksplisit dalam job YAML. Konfigurasi YAML eksplisit memiliki prioritas lebih tinggi. Sebagai contoh, dalam sampel di atas, parameter fs.oss.endpoint menggunakan nilai dari job YAML, menggantikan nilai di paimon_dlf_catalog.
Contoh
Saat menggunakan Paimon sebagai sink ingesti data, rujuk contoh berikut untuk mengonfigurasi job berdasarkan jenis katalog Paimon Anda.
-
Contoh konfigurasi untuk menulis ke Object Storage Service (OSS) dengan katalog Paimon
filesystem:source: type: mysql name: MySQL Source hostname: ${secret_values.mysql.hostname} port: ${mysql.port} username: ${secret_values.mysql.username} password: ${secret_values.mysql.password} tables: ${mysql.source.table} server-id: 8601-8604 sink: type: paimon name: Paimon Sink catalog.properties.metastore: filesystem catalog.properties.warehouse: oss://default/test catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com catalog.properties.fs.oss.accessKeyId: xxxxxxxx catalog.properties.fs.oss.accessKeySecret: xxxxxxxxUntuk informasi tentang parameter yang diawali dengan
catalog.properties, lihat Buat Katalog Filesystem Paimon. -
Contoh konfigurasi untuk menulis ke Data Lake Formation (DLF) dengan katalog Paimon
rest:source: type: mysql name: MySQL Source hostname: ${secret_values.mysql.hostname} port: ${mysql.port} username: ${secret_values.mysql.username} password: ${secret_values.mysql.password} tables: ${mysql.source.table} server-id: 8601-8604 sink: type: paimon name: Paimon Sink catalog.properties.metastore: rest catalog.properties.uri: dlf_uri catalog.properties.warehouse: your_warehouse catalog.properties.token.provider: dlf # (Opsional) Aktifkan deletion vectors untuk meningkatkan performa baca. table.properties.deletion-vectors.enabled: trueUntuk informasi tentang parameter yang diawali dengan
catalog.properties, lihat Parameter Konfigurasi Katalog Flink CDC.
Perubahan Skema
Saat digunakan sebagai sink ingesti data, Paimon mendukung event perubahan skema berikut:
-
CREATE TABLE EVENT
-
ADD COLUMN EVENT
-
ALTER COLUMN TYPE EVENT (Mengubah tipe data kolom kunci primer tidak didukung.)
-
RENAME COLUMN EVENT
-
DROP COLUMN EVENT
-
TRUNCATE TABLE EVENT
-
DROP TABLE EVENT
Jika tabel Paimon downstream sudah ada, job menulis ke skema yang ada dan tidak mencoba membuat tabel lagi.
FAQ
-
Mengapa job Paimon gagal dengan "Heartbeat of TaskManager timed out"?
-
Mengapa job Paimon gagal dengan "Sink materializer must not be used with Paimon sink"?
-
Mengapa job Paimon gagal dengan "File deletion conflicts detected" atau "LSM conflicts detected"?
-
Mengapa job Paimon gagal dengan "File xxx not found, Possible causes"?
-
Apakah visibilitas data connector Paimon terkait dengan interval checkpoint?