Impor data secara batch ke Hologres melalui konektor Flink untuk ingest data yang sangat efisien dan berbeban rendah.
Latar Belakang
Hologres terintegrasi dengan Apache Flink untuk menyediakan kemampuan streaming data real-time. Untuk kasus penggunaan yang tidak sensitif terhadap waktu—seperti memuat data historis, memproses data offline, atau mengagregasi log—impor batch merupakan pendekatan yang direkomendasikan. Impor batch menulis volume data besar ke Hologres sekaligus, sehingga lebih efisien dan menghemat sumber daya. Pilih antara impor real-time dan batch berdasarkan kebutuhan bisnis dan sumber daya yang tersedia. Untuk informasi lebih lanjut tentang impor real-time, lihat Realtime Compute for Apache Flink.
Prasyarat
-
Anda telah membeli instans Hologres. Untuk informasi lebih lanjut, lihat Beli instans Hologres.
-
Anda telah menerapkan kluster Apache Flink versi 1.15 atau yang lebih baru. Untuk informasi lebih lanjut, lihat topik berikut:
-
Apache Flink: Deploy Flink.
-
Realtime Compute for Apache Flink: Aktifkan Realtime Compute for Apache Flink.
-
Impor batch menggunakan Realtime Compute for Apache Flink
-
Buat tabel hasil Hologres untuk menyimpan data yang diimpor dari Flink. Untuk petunjuknya, lihat Hubungkan ke HoloWeb dan jalankan kueri. Topik ini menggunakan tabel
test_sink_customersebagai contoh.-- Buat tabel hasil Hologres. CREATE TABLE test_sink_customer ( c_custkey BIGINT, c_name TEXT, c_address TEXT, c_nationkey INT, c_phone TEXT, c_acctbal NUMERIC(15,2), c_mktsegment TEXT, c_comment TEXT, "date" DATE ) WITH ( distribution_key="c_custkey,date", orientation="column" );CatatanNama bidang dan tipe data di tabel sumber Flink harus sesuai dengan yang ada di tabel hasil Hologres.
-
Masuk ke Konsol Realtime Compute for Apache Flink. Di halaman Deployments, klik Create Deployment. Konfigurasikan parameter penerapan dan klik Deploy. Untuk informasi lebih lanjut tentang parameter tersebut, lihat Deploy a JAR job.
Tabel berikut menjelaskan parameter utama.
Parameter
Deskripsi
Deployment Type
Pilih JAR.
Deployment Mode
Anda dapat memilih mode stream atau mode batch. Topik ini menggunakan mode batch sebagai contoh.
Engine Version
Untuk informasi lebih lanjut tentang versi engine, lihat Engine versions dan Lifecycle policies. Topik ini menggunakan versi
vvr-8.0.7-flink-1.17sebagai contoh.JAR URI
Unggah konektor Flink open source: hologres-connector-flink-repartition.jar.
CatatanAnda dapat menggunakan konektor Flink open source untuk mengimpor data secara batch ke Hologres. Untuk kode sumber konektor Flink, lihat repositori GitHub resmi Hologres.
Entry Point Class
Kelas titik masuk program. Kelas utama untuk konektor Flink adalah
com.alibaba.ververica.connectors.hologres.example.FlinkToHoloRePartitionExample.Entry Point Main Arguments
Berikan path ke file
repartition.sql. Di runtime Realtime Compute for Apache Flink, file dependensi tambahan disimpan di/flink/usrlib/. Oleh karena itu, argumen lengkapnya adalah--sqlFilePath="/flink/usrlib/repartition.sql".Additional Dependencies
Unggah file
repartition.sql. File ini merupakan skrip Flink SQL yang digunakan untuk mendefinisikan sumber data, mendeklarasikan tabel hasil, dan mengonfigurasi koneksi ke Hologres. Kode berikut adalah contoh filerepartition.sql.-- DDL untuk tabel sumber. Contoh ini menggunakan konektor Flink DataGen untuk menghasilkan data uji. CREATE TEMPORARY TABLE source_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ) WITH ( 'connector' = 'datagen' ,'rows-per-second' = '10000' ,'number-of-rows' = '1000000' ); -- DQL untuk tabel sumber. Hasil kueri harus sesuai dengan skema tabel hasil yang didefinisikan dalam DDL sink, termasuk jumlah dan tipe bidang. SELECT *, cast('2024-04-21' as DATE) FROM source_table; -- DDL untuk tabel sink. Ini mendeklarasikan tabel hasil dan mengonfigurasi koneksi ke Hologres. CREATE TABLE sink_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ,`date` DATE ) WITH ( 'connector' = 'hologres' ,'dbname' = 'doc_****' ,'tablename' = 'test_sink_customer' ,'username' = 'yourAccessKeyId' ,'password' = 'yourAccessKeySecret' ,'endpoint' = 'hgpostcn-cn-7pp2e1k7****-cn-hangzhou.hologres.aliyuncs.com:80' ,'jdbccopywritemode' = 'true' ,'bulkload' = 'true' ,'target-shards.enabled'='true' );CatatanUntuk informasi lebih lanjut tentang parameter koneksi Hologres dalam file
repartition.sql, lihat Hologres Flink connector parameters. -
Klik nama penerapan dan buka halaman Deployment Details. Di bagian Resource Configurations, ubah Parallelism.
CatatanKami merekomendasikan agar Anda mengatur parallelism sesuai dengan ShardCount tabel hasil Hologres.
-
Kueri tabel hasil Hologres.
Setelah pekerjaan Flink dikirimkan, Anda dapat mengkueri data yang ditulis di Hologres. Contoh pernyataan:
SELECT * FROM test_sink_customer;
Impor batch menggunakan Apache Flink
-
Buat tabel hasil Hologres untuk menyimpan data yang diimpor dari Flink. Untuk petunjuknya, lihat Hubungkan ke HoloWeb dan jalankan kueri. Topik ini menggunakan tabel
test_sink_customersebagai contoh.-- Buat tabel hasil Hologres. CREATE TABLE test_sink_customer ( c_custkey BIGINT, c_name TEXT, c_address TEXT, c_nationkey INT, c_phone TEXT, c_acctbal NUMERIC(15,2), c_mktsegment TEXT, c_comment TEXT, "date" DATE ) WITH ( distribution_key="c_custkey,date", orientation="column" );CatatanAnda dapat mengatur jumlah shard berdasarkan volume data Anda. Untuk informasi lebih lanjut tentang shard, lihat Manage table groups and shard count.
-
Buat file
repartition.sqldan unggah ke direktori apa pun di lingkungan kluster Flink Anda. Topik ini menggunakan path/flink-1.15.4/src/repartition.sqlsebagai contoh. Kode berikut adalah contoh filerepartition.sql.CatatanFile ini merupakan skrip Flink SQL yang digunakan untuk mendefinisikan sumber data, mendeklarasikan tabel hasil, dan mengonfigurasi koneksi ke Hologres.
-- DDL untuk tabel sumber. Contoh ini menggunakan konektor Flink DataGen untuk menghasilkan data uji. CREATE TEMPORARY TABLE source_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ) WITH ( 'connector' = 'datagen' ,'rows-per-second' = '10000' ,'number-of-rows' = '1000000' ); -- DQL untuk tabel sumber. Hasil kueri harus sesuai dengan skema tabel hasil yang didefinisikan dalam DDL sink, termasuk jumlah dan tipe bidang. SELECT *, cast('2024-04-21' as DATE) FROM source_table; -- DDL untuk tabel sink. Ini mendeklarasikan tabel hasil dan mengonfigurasi koneksi ke Hologres. CREATE TABLE sink_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ,`date` DATE ) WITH ( 'connector' = 'hologres' ,'dbname' = 'doc_****' ,'tablename' = 'test_sink_customer' ,'username' = 'yourAccessKeyId' ,'password' = 'yourAccessKeySecret' ,'endpoint' = 'hgpostcn-cn-7pp2e1k7****-cn-hangzhou.hologres.aliyuncs.com:80' ,'jdbccopywritemode' = 'true' ,'bulkload' = 'true' ,'target-shards.enabled'='true' );Tabel berikut menjelaskan parameter utama.
Parameter
Wajib
Deskripsi
connector
Ya
Jenis konektor. Nilainya harus
hologres.dbname
Ya
Nama database Hologres.
tablename
Ya
Nama tabel Hologres yang menerima data.
username
Ya
ID AccessKey Akun Alibaba Cloud Anda.
Anda dapat memperoleh ID AccessKey dari halaman AccessKey Pair.
password
Ya
Rahasia AccessKey yang sesuai dengan ID AccessKey Anda.
endpoint
Ya
Titik akhir VPC instans Hologres. Buka halaman detail instans di Konsol Hologres dan peroleh titik akhir dari bagian Configurations.
CatatanTitik akhir harus mencantumkan nomor port dalam format
ip:port. Gunakan titik akhir VPC untuk koneksi dalam wilayah yang sama. Gunakan titik akhir publik untuk koneksi cross-region.jdbccopywritemode
Tidak
Metode penulisan data. Nilai yang valid:
-
false(default): Menggunakan metodeINSERT. -
true: Menggunakan metodeCOPY. MetodeCOPYmencakup streamingCOPY(Fixed Copy) dan batchCOPY. Secara default, streamingCOPYdigunakan.CatatanDibandingkan dengan metode
INSERT, streamingCOPYmenggunakan model streaming untuk mencapai throughput lebih tinggi, latensi data lebih rendah, dan konsumsi memori klien berkurang karena data tidak dibatch. Namun, metode ini tidak mendukung retraction data.
bulkload
Tidak
Apakah akan menggunakan metode batch
COPY. Nilai yang valid:-
true: Menggunakan batchCOPY. Pengaturan ini hanya berlaku jikajdbccopywritemodejuga diatur ketrue. Jika tidak, streamingCOPYyang digunakan.Catatan-
Dibandingkan dengan streaming
COPY, batchCOPYlebih efisien dan menggunakan sumber daya Hologres secara lebih efektif, sehingga menghasilkan performa penulisan yang lebih baik. Pilih metode penulisan yang sesuai berdasarkan kebutuhan bisnis Anda. -
Saat Anda menggunakan batch
COPYuntuk menulis data ke tabel dengan primary key, kunci tabel (table locks) mungkin terjadi. Anda dapat mengatur parametertarget-shards.enabledketrueuntuk mengurangi granularitas kunci dari level tabel ke level shard. Hal ini memungkinkan beberapa tugas impor batch berjalan secara konkuren dan mengurangi kontensi kunci tabel. Dibandingkan dengan streamingCOPY, pendekatan ini secara signifikan mengurangi beban pada instans Hologres saat menulis ke tabel dengan primary key. Pengujian menunjukkan pengurangan beban sekitar 66,7%. -
Saat Anda menggunakan batch
COPY, jika tabel tujuan memiliki primary key, tabel tersebut harus kosong sebelum operasi penulisan. Jika tidak, proses penulisan akan melambat karena deduplikasi data berdasarkan primary key.
-
-
false(default): Tidak menggunakan batchCOPY.
target-shards.enabled
Tidak
Apakah akan mengaktifkan penulisan batch ke shard target. Nilai yang valid:
-
true: Mengaktifkan penulisan batch ke shard target. Saat data sumber dipartisi ulang berdasarkan shard, hal ini mengurangi granularitas kunci ke level shard. -
false(default): Menonaktifkan fitur ini.
CatatanUntuk informasi lebih lanjut tentang parameter koneksi Hologres dalam file
repartition.sql, lihat Hologres Flink connector parameters. -
-
Di lingkungan kluster Flink Anda, unggah konektor Flink open source hologres-connector-flink-repartition.jar ke direktori apa pun. Topik ini menggunakan direktori root sebagai contoh.
CatatanAnda dapat menggunakan konektor Flink open source untuk mengimpor data secara batch ke Hologres. Untuk kode sumber konektor Flink, lihat repositori GitHub resmi Hologres.
-
Kirimkan pekerjaan Flink. Contoh perintah:
./bin/flink run -Dexecution.runtime-mode=BATCH -p 3 -c com.alibaba.ververica.connectors.hologres.example.FlinkToHoloRePartitionExample hologres-connector-flink-repartition.jar --sqlFilePath="/flink-1.15.4/src/repartition.sql"Parameter dalam perintah di atas:
-
Dexecution.runtime-mode: Mode eksekusi pekerjaan Flink. Untuk informasi lebih lanjut, lihat Execution Mode. -
p: Parallelisme pekerjaan. Kami merekomendasikan mengatur nilai ini sesuai dengan ShardCount tabel hasil, atau pembagi dari ShardCount tersebut. -
c: Nama lengkap kelas utama dalam file hologres-connector-flink-repartition.jar. -
sqlFilePath: Path ke filerepartition.sql.
-
-
Kueri tabel hasil Hologres.
Setelah pekerjaan Flink dikirimkan, Anda dapat mengkueri data yang ditulis di Hologres. Contoh pernyataan:
SELECT * FROM test_sink_customer;