Realtime Compute for Apache Flink menyederhanakan ingesti data real-time dengan secara otomatis menangani transisi dari sinkronisasi penuh ke inkremental, penemuan metadata, evolusi skema, dan sinkronisasi seluruh database. Topik ini menjelaskan cara membangun pekerjaan ingesti data yang mengalirkan data dari ApsaraDB RDS for MySQL ke Hologres.
Latar Belakang
Berikut ini menunjukkan struktur database dan tabel di konsol DMS.
Instans RDS berisi beberapa database (tpc_ds, tpc_ds_large, user_db1, user_db2, user_db3, __recycle_bin__). Database user_db3 berisi tiga tabel: user03, user06, dan user09. Tabel user03 memiliki dua kolom: id (int(11)) dan name (varchar(255)).
Ikuti langkah-langkah berikut untuk mengembangkan pekerjaan ingesti data yang menyinkronkan semua tabel tersebut ke Hologres dan menggabungkan tabel user terpartisi menjadi satu tabel:
Topik ini menggunakan ingesti data Flink CDC untuk melakukan sinkronisasi seluruh database dan menggabungkan tabel terpartisi. Pendekatan ini memungkinkan Anda menyelesaikan sinkronisasi data penuh dan inkremental, serta sinkronisasi perubahan skema real-time, dalam satu pekerjaan.
Prasyarat
-
Jika Anda menggunakan RAM user atau RAM role, pastikan Anda memiliki izin yang diperlukan untuk mengakses konsol Flink. Untuk informasi lebih lanjut, lihat Manage permissions.
-
Anda telah membuat ruang kerja Flink. Untuk informasi lebih lanjut, lihat Aktifkan Realtime Compute for Apache Flink.
-
Penyimpanan hulu dan hilir
-
Anda telah membuat instans ApsaraDB RDS for MySQL. Untuk informasi lebih lanjut, lihat (Deprecated, redirects to "Step 1") Buat instans ApsaraDB RDS for MySQL dengan cepat.
-
Anda telah membuat instans Hologres. Untuk informasi lebih lanjut, lihat Beli instans Hologres.
CatatanInstans ApsaraDB RDS for MySQL dan Hologres harus berada di wilayah dan virtual private cloud (VPC) yang sama dengan ruang kerja Flink. Jika tidak, Anda harus membuat koneksi jaringan. Untuk informasi lebih lanjut, lihat Bagaimana cara mengakses layanan lain lintas VPC? dan Bagaimana cara mengakses internet?.
-
-
Anda telah menyiapkan data uji dan mengonfigurasi daftar putih IP. Untuk informasi lebih lanjut, lihat Siapkan data uji MySQL dan database Hologres dan Konfigurasikan daftar putih IP.
Siapkan data uji MySQL dan database Hologres
-
Klik tpc_ds.sql, user_db1.sql, user_db2.sql, dan user_db3.sql untuk mengunduh file data uji ke mesin lokal Anda.
-
Di konsol DMS, siapkan data uji di instans ApsaraDB RDS for MySQL Anda.
-
Masuk ke instans ApsaraDB RDS for MySQL Anda menggunakan DMS.
Untuk informasi lebih lanjut, lihat (Deprecated, redirects to "Step 2") Masuk ke instans ApsaraDB RDS for MySQL menggunakan DMS.
-
Di jendela SQL Console, masukkan perintah berikut dan klik Execute.
Perintah berikut membuat empat database:
tpc_ds,user_db1,user_db2, danuser_db3.CREATE DATABASE tpc_ds; CREATE DATABASE user_db1; CREATE DATABASE user_db2; CREATE DATABASE user_db3; -
Di bilah navigasi atas, klik Data Import.
-
Di tab Batch Data Import, pilih database tempat mengimpor data, unggah file SQL yang sesuai, klik Submit, lalu klik Execute Change. Di kotak dialog yang muncul, klik Confirm Execution.
Ulangi operasi ini untuk mengimpor file data yang sesuai ke database
tpc_ds,user_db1,user_db2, danuser_db3.Atur File Encoding ke Auto Detect, Import Mode ke Express Mode atau Safe Mode, dan File Type ke SQL Script, CSV, atau Excel. Lampiran mendukung format txt/sql/csv/xlsx/zip, hingga 5 GB.
-
-
Di konsol Hologres, buat database bernama
my_useruntuk menyimpan data tabel user yang digabung.Untuk informasi lebih lanjut, lihat Buat database.
Konfigurasikan daftar putih IP
Untuk mengizinkan Flink mengakses instans ApsaraDB RDS for MySQL dan Hologres, tambahkan Blok CIDR ruang kerja Flink ke daftar putih IP kedua instans tersebut.
-
Dapatkan Blok CIDR ruang kerja Flink.
-
Masuk ke konsol Realtime Compute for Apache Flink.
-
Di daftar ruang kerja, temukan workspace target dan pilih di kolom Actions.
-
Di kotak dialog Workspace Details, lihat informasi CIDR Block untuk vSwitch Flink.
-
-
Tambahkan Blok CIDR Flink ke daftar putih IP instans ApsaraDB RDS for MySQL.
Untuk informasi lebih lanjut, lihat Konfigurasikan daftar putih IP.
Di kotak dialog Modify Whitelist Group, masukkan Blok CIDR Flink fully managed di bidang Group Whitelist, lalu klik OK.
-
Tambahkan Blok CIDR Flink ke daftar putih IP instans Hologres.
Saat mengonfigurasi koneksi data di HoloWeb, atur Login Method ke Password-free login for current user sebelum mengonfigurasi daftar putih IP untuk koneksi tersebut. Untuk informasi lebih lanjut, lihat daftar putih IP.
Di HoloWeb, pilih Security Center dari bilah navigasi atas, klik IP Whitelist di panel kiri, dan konfigurasikan parameter berikut di kotak dialog Edit IP Whitelist:
-
Group: Pilih default
-
Database Restriction: Pilih ALL
-
User Restriction: Pilih ALL
-
IP Address: Masukkan Blok CIDR Flink dalam Notasi CIDR (misalnya, 172.xx.0/19). Pisahkan beberapa alamat IP dengan baris baru.
Klik OK untuk menerapkan.
-
Langkah 1: Kembangkan pekerjaan ingesti
-
Masuk ke konsol pengembangan Flink dan buat pekerjaan baru.
-
Di halaman , klik New.
-
Klik Blank Data Ingestion Draft.
Realtime Compute for Apache Flink menyediakan berbagai templat kode, masing-masing dengan kasus penggunaan spesifik, contoh kode, dan panduan. Anda dapat mengklik templat untuk mempelajari fitur produk dan sintaks guna mengimplementasikan logika bisnis Anda.
-
Klik Next.
-
Di kotak dialog New Draft, konfigurasikan parameter.
Parameter
Deskripsi
Contoh
Name
Nama pekerjaan.
CatatanNama pekerjaan harus unik dalam proyek saat ini.
flink-test
Location
Folder tempat file kode pekerjaan disimpan.
Anda juga dapat mengklik ikon
di samping folder yang ada untuk membuat subfolder.Job Drafts
Engine Version
Versi engine Flink yang digunakan oleh pekerjaan. Untuk informasi tentang nomor versi engine, kompatibilitas versi, dan tanggal penting siklus hidup, lihat Engine Versions.
vvr-11.1-jdk11-flink-1.20
-
Klik OK.
-
-
Salin kode pekerjaan berikut ke editor pekerjaan.
Kode berikut menyinkronkan semua tabel dari database
tpc_dske Hologres dan menggabungkan tabel user terpartisi menjadi satu tabel di Hologres:source: type: mysql name: MySQL Source hostname: localhost port: 3306 username: username password: password tables: tpc_ds.\.*,user_db[0-9]+.user[0-9]+ server-id: 8601-8604 # (Optional) Synchronize table and column comments. include-comments.enabled: true # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory errors. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to accelerate reading. scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: hologres name: Hologres Sink endpoint: ****.hologres.aliyuncs.com:80 dbname: cdcyaml_test username: ${secret_values.holo-username} password: ${secret_values.holo-password} sink.type-normalize-strategy: BROADEN route: # Merge and synchronize the sharded user tables to the my_user.users table. - source-table: user_db[0-9]+.user[0-9]+ sink-table: my_user.usersCatatanTabel dari database MySQL
tpc_dsdipetakan langsung ke tabel dengan nama yang sama di tujuan, sehingga bagianroutetidak memerlukan konfigurasi pemetaan tambahan. Untuk menyinkronkan tabel ke database dengan nama berbeda, sepertiods_tps_ds, konfigurasikan modulroutesebagai berikut:route: # Merge and synchronize the sharded user tables to the my_user.users table. - source-table: user_db[0-9]+.user[0-9]+ sink-table: my_user.users # Rename the database for all tables under tpc_ds and synchronize them to ods_tps_ds. - source-table: tpc_ds.\.* sink-table: ods_tps_ds.<> replace-symbol: <>
Langkah 2: Mulai pekerjaan
-
Di halaman , klik Deploy. Di kotak dialog yang muncul, klik Confirm.
Di kotak dialog, Anda dapat memasukkan Remarks, mengatur Job Tags, dan memilih antrian dari daftar drop-down Deployment Target (default:
default-queue). Catatan: Penerapan berlaku saat pekerjaan dimulai berikutnya. -
Di halaman , klik Start di kolom Actions untuk pekerjaan target. Konfigurasikan parameter sesuai kebutuhan. Untuk informasi lebih lanjut, lihat Mulai pekerjaan.
-
Klik Start.
Setelah pekerjaan dimulai, Anda dapat melihat informasi waktu proses dan statusnya di halaman Job Operations. Status pekerjaan mencakup Running, Failed, dan Stopped. Gunakan filter drop-down Stream Jobs di bagian atas untuk memfilter daftar pekerjaan berdasarkan jenis.
Langkah 3: Verifikasi sinkronisasi penuh
Masuk ke Hologres Management Console.
-
Di tab Metadata Management, verifikasi bahwa 24 tabel dan datanya tersedia di database
tpc_dsinstans Hologres.Navigasi ke Instance > tpc_ds > public > Tables untuk melihat tabel seperti call_center, catalog_sales, customer, dan store_sales. Pilih tabel store_sales dan klik tab Data Preview untuk melihat kolom seperti ss_sold_date, ss_sold_time, ss_item_sk, dan ss_customer.
-
Di tab Metadata Management, periksa skema tabel
usersdi databasemy_user.Berikut ini menjelaskan skema dan data tabel yang disinkronkan.
-
Skema tabel
Setelah sinkronisasi, tabel
usersberisi kolom-kolom berikut:-
_db_name(text, primary key) -
_table_name(text, primary key) -
id(int4, primary key) -
name(varchar 255, nullable)
Kolom
_db_name,_table_name, danidmembentuk kunci utama gabungan.Skema tabel
usersmencakup dua kolom tambahan yang tidak ada di tabel sumber MySQL:_db_namedan_table_name. Kolom-kolom ini menunjukkan database dan tabel sumber untuk setiap baris dan merupakan bagian dari kunci utama gabungan, memastikan keunikan data setelah tabel terpartisi digabung. -
-
Data tabel
Di pojok kanan atas halaman informasi tabel
users, klik Query Table. Masukkan perintah berikut, lalu klik Run.select * from users order by _db_name,_table_name,id;Hasil kueri menunjukkan bahwa data user telah sepenuhnya disinkronkan ke tiga database:
user_db1(user01, user04, user07),user_db2(user02, user05, user08), danuser_db3(user03, user06, user09), total 9 catatan dengan kolom_db_name,_table_name,id, danname.
-
Langkah 4: Verifikasi sinkronisasi inkremental
Setelah sinkronisasi penuh selesai, pekerjaan secara otomatis beralih ke fase sinkronisasi inkremental tanpa intervensi manual. Anda dapat memeriksa nilai currentEmitEventTimeLag di tab Monitoring and Alerts untuk menentukan fase sinkronisasi data.
-
Masuk ke konsol Realtime Compute for Apache Flink.
-
Klik Console di kolom Actions untuk ruang kerja target.
-
Di halaman , klik nama pekerjaan target.
-
Klik tab Alarm (atau Metrics).
-
Periksa grafik
currentEmitEventTimeLaguntuk menentukan fase sinkronisasi data.
-
Nilai 0 menunjukkan fase sinkronisasi penuh.
-
Nilai lebih besar dari 0 menunjukkan fase sinkronisasi inkremental.
-
-
Verifikasi sinkronisasi data dan perubahan skema real-time.
Sumber CDC MySQL mendukung sinkronisasi data dan skema real-time selama fase inkremental. Untuk memverifikasi hal ini, ubah skema dan data tabel user terpartisi di MySQL setelah pekerjaan memasuki fase ini.
-
Masuk ke instans ApsaraDB RDS for MySQL menggunakan DMS.
Untuk informasi lebih lanjut, lihat (Deprecated, redirects to "Step 2") Masuk ke instans ApsaraDB RDS for MySQL menggunakan DMS.
-
Di database
user_db2, jalankan perintah berikut untuk mengubah skema tabeluser02serta menyisipkan dan memperbarui data.USE `user_db2`; ALTER TABLE `user02` ADD COLUMN `age` INT; -- Add the age column. INSERT INTO `user02` (id, name, age) VALUES (27, 'Tony', 30); -- Insert a row that includes the age data. UPDATE `user05` SET name='JARK' WHERE id=15; -- Update another table and change the name to uppercase. -
Di konsol Hologres, periksa perubahan skema dan data tabel
users.Di pojok kanan atas halaman informasi tabel
users, klik Query Table, masukkan perintah berikut, lalu klik Run.select * from users order by _db_name,_table_name,id;Perubahan skema pada
user02dan modifikasi data disebarkan secara real-time, meskipun tabel terpartisi memiliki skema berbeda. TabelusersHologres kini menampilkan kolomagebaru, catatan yang disisipkan untuk Tony, dan catatan yang diperbarui untuk JARK.Hasil kueri menunjukkan 11 baris dari tabel users di user_db1, user_db2, dan user_db3. Perubahan sinkronisasi inkremental tercermin: user Tony (id 27) di user_db2 memiliki
agediperbarui menjadi 30, dan catatan lain memilikinamediubah menjadi JARK. Baris-baris lain memiliki nilai age \N, yang mengonfirmasi bahwa sinkronisasi inkremental berfungsi.
-
(Opsional) Langkah 5: Konfigurasikan resource pekerjaan
Untuk kinerja yang lebih baik, Anda dapat menyesuaikan resource pekerjaan seperti konkurensi, memori TaskManager, dan CUs berdasarkan volume data Anda.
-
Di halaman , klik nama pekerjaan target.
-
Di tab Configuration, klik Edit di pojok kanan atas bagian Resources.
-
Atur manual parameter resource seperti memori TaskManager dan konkurensi.
-
Di sisi kanan bagian Resources, klik Save.
-
Restart pekerjaan.
Perubahan konfigurasi resource hanya berlaku setelah Anda me-restart pekerjaan.
Dokumen terkait
-
Untuk sintaks setiap modul ingesti data, lihat Referensi pengembangan pekerjaan ingesti data Flink CDC.
-
Jika Anda mengalami masalah saat pekerjaan ingesti data berjalan, lihat Masalah umum dan solusi untuk pekerjaan ingesti data.