E-MapReduce mendukung pembacaan dan penulisan data ke Paimon menggunakan Flink SQL. Topik ini menyediakan contoh cara membuat katalog, melakukan pembacaan dan penulisan streaming, serta menjalankan kueri OLAP.
Prasyarat
Anda telah membuat kluster Dataflow atau kluster kustom dengan Flink dan Paimon yang dipilih. Untuk informasi selengkapnya, lihat Buat kluster.
Untuk menggunakan katalog Hive, Anda harus membuat kluster kustom dengan Flink, Paimon, dan Hive yang dipilih. Anda juga harus mengatur tipe Metadata Storage Method menjadi Self-managed RDS atau Built-in MySQL.
Batasan
-
E-MapReduce V3.46.0 dan V5.17.0 tidak mendukung katalog DLF dan katalog Hive.
-
Anda dapat menggunakan Flink SQL untuk membaca dari dan menulis ke Paimon pada kluster yang menjalankan E-MapReduce V3.46.0 hingga V3.50.X dan E-MapReduce V5.12.0 hingga V5.16.X.
CatatanUntuk E-MapReduce V3.51.X dan versi lebih baru serta E-MapReduce V5.17.X dan versi lebih baru, rujuk ke dokumentasi Apache Paimon dan konfigurasikan integrasi tersebut di kluster EMR Anda.
Prosedur
Langkah 1: Konfigurasikan dependensi
Anda dapat membaca dari dan menulis ke Paimon menggunakan katalog Filesystem, katalog Hive, atau katalog DLF. Konfigurasikan dependensi sesuai metode yang Anda pilih.
Katalog Filesystem
cp /opt/apps/PAIMON/paimon-current/lib/flink/*.jar /opt/apps/FLINK/flink-current/lib/
Katalog Hive
cp /opt/apps/PAIMON/paimon-current/lib/flink/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/FLINK/flink-current/opt/catalogs/hive-2.3.6/*.jar /opt/apps/FLINK/flink-current/lib/
Katalog DLF
cp /opt/apps/PAIMON/paimon-current/lib/flink/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/PAIMON/paimon-current/lib/jackson/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/METASTORE/metastore-*/hive2/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/FLINK/flink-current/opt/catalogs/hive-2.3.6/*.jar /opt/apps/FLINK/flink-current/lib/
Langkah 2: Mulai kluster
Topik ini menggunakan mode session sebagai contoh. Untuk mode lainnya, lihat Penggunaan dasar.
Jalankan perintah berikut untuk memulai session YARN terpisah (detached):
yarn-session.sh --detached
Langkah 3: Buat katalog
Paimon menyimpan data dan metadata di sistem file seperti HDFS atau di penyimpanan objek seperti OSS-HDFS. Parameter warehouse menentukan path root. Jika path warehouse yang ditentukan belum ada, Paimon akan membuatnya secara otomatis. Jika path tersebut sudah ada, Anda dapat menggunakan katalog tersebut untuk mengakses tabel-tabel yang sudah ada di path tersebut.
Anda juga dapat menyinkronkan metadata ke Hive atau DLF agar layanan lain dapat mengakses data Paimon.
E-MapReduce V3.46.0 dan V5.17.0 tidak mendukung katalog DLF dan katalog Hive.
Katalog Filesystem
Katalog Filesystem hanya menyimpan metadata di sistem file atau penyimpanan objek.
-
Jalankan perintah berikut untuk memulai client Flink SQL.
sql-client.sh -
Jalankan pernyataan Flink SQL berikut untuk membuat katalog Filesystem.
CREATE CATALOG test_catalog WITH ( 'type' = 'paimon', 'metastore' = 'filesystem', 'warehouse' = 'oss://<yourBucketName>/warehouse' );
Katalog Hive
Katalog Hive menyinkronkan metadata ke Hive Metastore. Tabel yang dibuat dalam katalog Hive dapat langsung dikueri dari Hive.
Untuk detail tentang kueri Paimon dari Hive, lihat Integrasikan Paimon dengan Hive.
-
Jalankan perintah berikut untuk memulai client Flink SQL.
sql-client.shCatatanPerintah startup sama terlepas dari versi Hive Anda.
-
Jalankan pernyataan Flink SQL berikut untuk membuat katalog Hive.
CREATE CATALOG test_catalog WITH ( 'type' = 'paimon', 'metastore' = 'hive', 'uri' = 'thrift://master-1-1:9083', -- Parameter uri menentukan alamat layanan Hive metastore. 'warehouse' = 'oss://<yourBucketName>/warehouse' );
Katalog DLF
Katalog DLF menyinkronkan metadata ke DLF.
Saat membuat kluster, Metadata harus diatur ke DLF Unified Metadata.
-
Jalankan perintah berikut untuk memulai client Flink SQL.
sql-client.shCatatanPerintah startup sama terlepas dari versi Hive Anda.
-
Jalankan pernyataan Flink SQL berikut untuk membuat katalog DLF.
CREATE CATALOG test_catalog WITH ( 'type' = 'paimon', 'metastore' = 'dlf', 'hive-conf-dir' = '/etc/taihao-apps/flink-conf', 'warehouse' = 'oss://<yourBucketName>/warehouse' );
Langkah 4: Baca dan tulis Paimon melalui streaming
Jalankan pernyataan Flink SQL berikut untuk membuat tabel di katalog, lalu baca dari dan tulis ke tabel tersebut.
-- Atur mode eksekusi ke streaming.
SET 'execution.runtime-mode' = 'streaming';
-- Paimon mengharuskan Anda mengatur interval checkpoint untuk pekerjaan streaming.
SET 'execution.checkpointing.interval' = '10s';
-- Gunakan katalog yang dibuat pada langkah sebelumnya.
USE CATALOG test_catalog;
-- Buat dan gunakan database uji coba.
CREATE DATABASE test_db;
USE test_db;
-- Gunakan datagen untuk menghasilkan data acak.
CREATE TEMPORARY TABLE datagen_source (
uuid int,
kind int,
price int
) WITH (
'connector' = 'datagen',
'fields.kind.min' = '0',
'fields.kind.max' = '9',
'rows-per-second' = '10'
);
-- Buat tabel Paimon.
CREATE TABLE test_tbl (
uuid int,
kind int,
price int,
PRIMARY KEY (uuid) NOT ENFORCED
);
-- Tulis data ke tabel Paimon.
INSERT INTO test_tbl SELECT * FROM datagen_source;
-- Baca data dari tabel.
-- Pekerjaan penulisan streaming sebelumnya berjalan secara konkuren.
-- Pastikan kluster Flink Anda memiliki sumber daya yang cukup (slot tugas) untuk kedua pekerjaan tersebut. Jika tidak, kueri ini tidak akan dieksekusi.
SELECT kind, SUM(price) FROM test_tbl GROUP BY kind;
Langkah 5: Jalankan kueri OLAP pada Paimon
Jalankan pernyataan Flink SQL berikut untuk melakukan kueri OLAP pada tabel yang telah dibuat.
-- Atur mode eksekusi ke batch.
RESET 'execution.checkpointing.interval';
SET 'execution.runtime-mode' = 'batch';
-- Gunakan mode tableau untuk mencetak hasil langsung di terminal.
SET 'sql-client.execution.result-mode' = 'tableau';
-- Kueri data di tabel.
SELECT kind, SUM(price) FROM test_tbl GROUP BY kind;
Langkah 6: Bersihkan resource
Setelah selesai pengujian, hentikan pekerjaan penulisan streaming Paimon untuk mencegah kebocoran resource.
Setelah pekerjaan dihentikan, jalankan pernyataan Flink SQL berikut untuk menghapus tabel yang telah dibuat.
DROP TABLE test_tbl;