Pelajari cara mengakses katalog Data Lake Formation (DLF) menggunakan Paimon REST di lingkungan EMR on ECS Spark.
Prasyarat
Buat kluster EMR versi 5.12.0 atau yang lebih baru dan pilih Spark 3 serta Paimon sebagai komponen. Untuk informasi tentang persyaratan versi lainnya, bergabunglah dengan grup DingTalk (106575000021) untuk menghubungi pengembang DLF.
Anda telah menyelesaikan Panduan Cepat untuk DLF.
Pastikan EMR dan DLF berada di wilayah yang sama, dan tambahkan VPC kluster EMR Anda ke daftar putih DLF.
Buat katalog
Lihat Menyiapkan DLF.
Berikan Izin DLF kepada Role
Berikan izin RAM kepada role AliyunECSInstanceForEMRRole. Langkah ini tidak diperlukan setelah produk EMR terintegrasi.
Masuk ke Konsol Resource Access Management (RAM) menggunakan Akun Alibaba Cloud Anda atau sebagai administrator RAM.
Di panel navigasi sebelah kiri, pilih dan cari role AliyunECSInstanceForEMRRole.
Di kolom Actions, klik Add Permissions untuk membuka halaman Add Permissions.
Di bawah Permission Policies, cari dan pilih AliyunDLFFullAccess, lalu klik Confirm.

Berikan izin DLF kepada role AliyunECSInstanceForEMRRole.
Masuk ke Konsol Data Lake Formation.
Di halaman Catalogs, klik nama katalog untuk membuka halaman detailnya.
Untuk memberikan izin ke seluruh katalog, klik tab Permissions. Atau, Anda dapat menavigasi ke database atau tabel tertentu dan mengklik tab Permissions-nya untuk memberikan akses.
Di halaman otorisasi, konfigurasikan pengaturan berikut dan klik OK.
User/Role: Pilih RAM User/RAM Role.
Select Authorization Object: Pilih AliyunECSInstanceForEMRRole dari daftar dropdown.
CatatanJika AliyunECSInstanceForEMRRole tidak muncul dalam daftar dropdown, buka halaman manajemen pengguna dan klik Sync.
Preset Permission Type: Pilih izin baca secara manual atau gunakan role yang telah ditentukan sebelumnya, seperti Data Reader atau Data Editor.
Upgrade Dependensi Paimon di Kluster EMR Anda
Unduh dua file JAR untuk Paimon versi 1.1 atau yang lebih baru dari Repositori Maven: paimon-jindo-*.jar dan paimon-spark-3.x-*.jar. Pastikan versi yang dipilih sesuai dengan versi Spark di kluster EMR Anda.
Impor dependensi Paimon.
Unggah dua file JAR,
paimon-jindo-*.jardanpaimon-spark-3.x-*.jar, ke OSS dan atur izin file-nya menjadi public-read. Untuk informasi selengkapnya, lihat Simple Upload.Ubah dan unggah skrip berikut ke OSS.
#!/bin/bash echo 'clean up paimon-dlf-2.5 exists file' rm -rf /opt/apps/PAIMON/paimon-dlf-2.5 rm -rf /opt/apps/PAIMON/paimon-dlf-2.5.tar.gz.* cd /opt/apps/PAIMON/paimon-current/lib/spark3 mkdir -p /opt/apps/PAIMON/paimon-dlf-2.5/lib/spark3 cd /opt/apps/PAIMON/paimon-dlf-2.5/lib/spark3 wget ${paimon-jindo-1.1.0.jar} wget ${paimon-spark-3.x-1.1.0.jar} echo 'link paimon-current to paimon-dlf-2.5' rm -f /opt/apps/PAIMON/paimon-current ln -sf /opt/apps/PAIMON/paimon-dlf-2.5 /opt/apps/PAIMON/paimon-currentPentingGanti placeholder
${paimon-jindo-1.1.0.jar}dan${paimon-spark-3.x-1.1.0.jar}dengan URL unduhan OSS yang sebenarnya. Secara default, kluster EMR on ECS tidak dapat mengakses jaringan publik.Jaringan privat:
https://{bucket}.oss-cn-hangzhou-internal.aliyuncs.com/jars/paimon-jindo-1.1.0.jar.Jaringan publik:
https://{bucket}.oss-cn-hangzhou.aliyuncs.com/jars/paimon-jindo-1.1.0.jar.
Jalankan skrip tersebut sebagai tindakan bootstrap pada kluster EMR Anda. Untuk informasi selengkapnya, lihat Manually Run a Script.
Di kluster EMR, pada tab , klik Create And Execute.
Di kotak dialog yang muncul, konfigurasikan pengaturan berikut dan klik OK.
Name: Masukkan nama skrip kustom.
Script Location: Pilih skrip upgrade yang telah Anda unggah ke OSS. Jalur harus dalam format oss://**/*.sh.
Execution Scope: Pilih Cluster.
Setelah skrip dijalankan, restart layanan Spark agar perubahan berlaku.
Baca dan Tulis Data dengan Spark
Hubungkan ke Katalog Paimon
Jalankan perintah spark-sql berikut di terminal.
Ganti ${regionID} dengan ID wilayah Anda, misalnya cn-hangzhou, dan ganti ${catalog} dengan nama katalog yang telah Anda buat di DLF.
spark-sql --master yarn \
--conf spark.driver.memory=5g \
--conf spark.sql.defaultCatalog=paimon \
--conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \
--conf spark.sql.catalog.paimon.metastore=rest \
--conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions \
--conf spark.sql.catalog.paimon.uri=http://${regionID}-vpc.dlf.aliyuncs.com \
--conf spark.sql.catalog.paimon.warehouse=${catalog} \
--conf spark.sql.catalog.paimon.token.provider=dlf \
--conf spark.sql.catalog.paimon.dlf.token-loader=ecsBuat Tabel Data
Jalankan pernyataan SQL berikut untuk membuat tabel data.
CREATE TABLE user_samples
(
user_id INT,
age INT,
gender_code STRING,
clk BOOLEAN
);
CREATE TABLE user_samples_di (
user_id INT,
age INT,
gender_code STRING,
clk BOOLEAN
)
USING CSV
OPTIONS(
'path'='oss://${bucket}/user/user_samples_di'
);Jika Anda tidak menentukan database, tabel akan dibuat di database
defaultkatalog. Anda juga dapat membuat dan menggunakan database lain.Anda harus membuat folder
/user/user_samplesdi OSS terlebih dahulu. Jika Anda menentukanpath, tabel eksternal akan dibuat. Dalam kasus ini, metadata disimpan dan dikelola di DLF, tetapi file data disimpan di jalur OSS yang ditentukan. Jika Anda menghapus tabel, hanya metadata yang dihapus. File data asli di OSS tidak dihapus.
Masukkan Data
Jalankan pernyataan SQL berikut untuk memasukkan data.
INSERT INTO user_samples VALUES
(1, 25, 'M', true),
(2, 18, 'F', false);
INSERT INTO user_samples_di VALUES
(1, 25, 'M', true),
(2, 18, 'F', true),
(3, 35, 'M', true);Kueri Data
Jalankan pernyataan SQL berikut untuk mengkueri data.
SELECT * FROM user_samples;
SELECT * FROM user_samples_di;Hasil berikut dikembalikan.


Gabungkan Data
Gabungkan data dari tabel user_samples_di ke tabel user_samples:
MERGE INTO user_samples
USING user_samples_di
ON user_samples.user_id = user_samples_di.user_id
WHEN MATCHED THEN
UPDATE SET
age = user_samples_di.age,
gender_code = user_samples_di.gender_code,
clk = user_samples_di.clk
WHEN NOT MATCHED THEN
INSERT (user_id, age, gender_code, clk)
VALUES (user_samples_di.user_id, user_samples_di.age, user_samples_di.gender_code, user_samples_di.clk);Setelah operasi merge, data di tabel user_samples ditimpa oleh data dari tabel user_samples_di untuk catatan yang memiliki user_id yang sama. Data baru dimasukkan ke tabel user_samples jika user_id yang sesuai tidak ditemukan.
