Realtime Compute for Apache Flink menyederhanakan ingesti data waktu nyata 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
Gambar berikut menunjukkan tampilan database dan tabel di Konsol DMS.
Ikuti langkah-langkah berikut untuk mengembangkan pekerjaan ingesti data yang menyinkronkan semua tabel tersebut ke Hologres serta menggabungkan tabel pengguna terbagi (sharded) menjadi satu tabel:
Topik ini menggunakan ingesti data Flink CDC untuk melakukan sinkronisasi seluruh database dan menggabungkan tabel terbagi. Pendekatan ini memungkinkan Anda menyelesaikan sinkronisasi data penuh dan inkremental, serta sinkronisasi perubahan skema waktu nyata, 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 Mengelola izin.
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 (Tidak digunakan lagi, dialihkan ke "Langkah 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 (Tidak digunakan lagi, dialihkan ke "Langkah 2") Masuk ke instans ApsaraDB RDS for MySQL menggunakan DMS.
Di jendela SQL Console, masukkan perintah berikut lalu 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 langkah ini untuk mengimpor file data yang sesuai ke database
tpc_ds,user_db1,user_db2, danuser_db3.
Di Konsol Hologres, buat database bernama
my_useruntuk menyimpan data tabel pengguna 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 lalu pilih di kolom Actions.
Di 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.

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.

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 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 dialog New Data Ingestion Job Draft, konfigurasikan parameter pekerjaan.
Parameter
Deskripsi
Contoh
File Name
Nama pekerjaan.
CatatanNama pekerjaan harus unik dalam proyek saat ini.
flink-test
Storage 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, kompatibilitas versi, dan tanggal penting siklus hidup, lihat Versi Engine.
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 pengguna terbagi 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 # (Opsional) Sinkronkan komentar tabel dan kolom. include-comments.enabled: true # (Opsional) Utamakan distribusi chunk tak terbatas untuk mencegah potensi error OutOfMemory pada TaskManager. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Opsional) Aktifkan filter parsing untuk mempercepat pembacaan. 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: # Gabungkan dan sinkronkan tabel pengguna terbagi ke tabel my_user.users. - 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: # Gabungkan dan sinkronkan tabel pengguna terbagi ke tabel my_user.users. - source-table: user_db[0-9]+.user[0-9]+ sink-table: my_user.users # Ubah nama database untuk semua tabel di bawah tpc_ds dan sinkronkan ke ods_tps_ds. - source-table: tpc_ds.\.* sink-table: ods_tps_ds.<> replace-symbol: <>
Langkah 2: Mulai pekerjaan
Di halaman , klik Deploy. Di dialog yang muncul, klik Confirm.

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.

Langkah 3: Verifikasi sinkronisasi penuh
Masuk ke Konsol Manajemen Hologres.
Di tab Metadata Management, verifikasi bahwa 24 tabel beserta datanya tersedia di database
tpc_dsinstans Hologres.
Di tab Metadata Management, periksa skema tabel
usersdi databasemy_user.Gambar berikut menunjukkan skema dan data tabel yang telah disinkronkan.
Skema tabel

Skema tabel
usersmencakup dua kolom tambahan yang tidak ada di tabel MySQL sumber:_db_namedan_table_name. Kolom-kolom ini menunjukkan database dan tabel sumber untuk setiap baris dan merupakan bagian dari kunci primer komposit, memastikan keunikan data setelah tabel terbagi 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 ditunjukkan pada gambar berikut.

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 Monitoring and Alerts (atau Metrics).
Periksa grafik
currentEmitEventTimeLaguntuk menentukan fase sinkronisasi data.
Nilai 0 menunjukkan fase sinkronisasi penuh.
Nilai lebih dari 0 menunjukkan fase sinkronisasi inkremental.
Verifikasi sinkronisasi data dan perubahan skema waktu nyata.
Sumber MySQL CDC mendukung sinkronisasi data dan skema waktu nyata selama fase inkremental. Untuk memverifikasi hal ini, ubah skema dan data tabel pengguna terbagi di MySQL setelah pekerjaan memasuki fase ini.
Masuk ke instans ApsaraDB RDS for MySQL menggunakan DMS.
Untuk informasi lebih lanjut, lihat (Tidak digunakan lagi, dialihkan ke "Langkah 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; -- Tambahkan kolom age. INSERT INTO `user02` (id, name, age) VALUES (27, 'Tony', 30); -- Sisipkan baris yang mencakup data age. UPDATE `user05` SET name='JARK' WHERE id=15; -- Perbarui tabel lain dan ubah nama menjadi huruf kapital.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;Gambar berikut menunjukkan hasil kueri. Perubahan skema pada
user02dan modifikasi data disebarkan secara waktu nyata, meskipun tabel terbagi memiliki skema berbeda. TabelusersHologres kini menampilkan kolomagebaru, catatan yang disisipkan untuk Tony, dan catatan yang diperbarui untuk JARK.
(Opsional) Langkah 5: Konfigurasikan sumber daya pekerjaan
Untuk performa yang lebih baik, Anda dapat menyesuaikan sumber daya pekerjaan seperti konkurensi, memori TaskManager, dan CUs berdasarkan volume data Anda.
Di halaman , klik nama pekerjaan target.
Di tab Deployment Details, klik Edit di pojok kanan atas bagian Resource Configuration.
Atur secara manual parameter sumber daya seperti memori TaskManager dan konkurensi.
Di sisi kanan bagian Resource Configuration, klik Save.
Restart pekerjaan.
Perubahan konfigurasi sumber daya 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.