Realtime Compute for Apache Flink memungkinkan Anda membuat pekerjaan ingesti data Flink CDC menggunakan file YAML untuk menyinkronkan data dari sumber ke sink. Topik ini menjelaskan langkah-langkah pengembangan pekerjaan ingesti data Flink CDC.
Latar Belakang
Dengan konfigurasi YAML, Anda dapat dengan mudah mendefinisikan pipa ETL kompleks yang secara otomatis dikonversi menjadi logika eksekusi Flink. Ingesti data Flink CDC menyediakan solusi integrasi data berbasis Flink CDC. Solusi ini secara efisien mendukung sinkronisasi seluruh database, sinkronisasi tabel tunggal, sinkronisasi database dan tabel terpartisi (sharded), penemuan otomatis tabel baru, penanganan perubahan skema, serta kolom terhitung kustom. Solusi ini juga mendukung pemrosesan ETL, penyaringan klausa WHERE, dan pemangkasan kolom. Pendekatan deklaratif ini secara signifikan menyederhanakan proses integrasi data serta meningkatkan efisiensi dan keandalannya.
Keunggulan Flink CDC
Di Realtime Compute for Apache Flink, Anda dapat mengembangkan pekerjaan ingesti data Flink CDC, pekerjaan SQL, atau pekerjaan DataStream untuk melakukan sinkronisasi data. Bagian berikut menjelaskan keunggulan penggunaan pekerjaan ingesti data Flink CDC dibandingkan dua opsi lainnya.
Flink CDC vs. Flink SQL
Pekerjaan ingesti data Flink CDC dan pekerjaan SQL menggunakan tipe data berbeda untuk transmisi data:
-
Pekerjaan SQL mentransmisikan data sebagai
RowData. Setiap objekRowDatamemiliki tipe perubahan sendiri, yang mencakup empat tipe utama: insert (+I), update before (-U), update after (+U), dan delete (-D). -
Flink CDC menggunakan
SchemaChangeEventuntuk menyampaikan informasi perubahan skema, seperti membuat tabel, menambah kolom, atau memotong (truncate) tabel. Flink CDC menggunakanDataChangeEventuntuk menyampaikan perubahan data, seperti insert, update, dan delete. Pesan update berisi konten sebelum dan sesudah perubahan, sehingga Anda dapat menulis data perubahan asli ke sink.
Tabel berikut mencantumkan keunggulan pekerjaan ingesti data Flink CDC dibandingkan pekerjaan SQL.
|
Flink CDC Ingestion |
Flink SQL |
|
Deteksi skema otomatis dan sinkronisasi seluruh database |
Memerlukan pernyataan manual |
|
Mendukung berbagai kebijakan untuk perubahan skema |
Tidak mendukung perubahan skema |
|
Mempertahankan changelog asli |
Mengganggu struktur changelog asli |
|
Mendukung pembacaan dari dan penulisan ke beberapa tabel |
Membaca dari dan menulis ke satu tabel saja |
Dibandingkan pernyataan CTAS/CDAS, pekerjaan Flink CDC menawarkan fitur yang lebih kuat, termasuk:
-
Sinkronisasi langsung perubahan skema hulu, tanpa menunggu penulisan data baru untuk memicu sinkronisasi.
-
Mempertahankan changelog asli, memastikan pesan update tidak dipisah.
-
Sinkronisasi lebih banyak jenis perubahan skema, seperti
TRUNCATE TABLEdanDROP TABLE. -
Pemetaan tabel fleksibel dan definisi nama tabel sink.
-
Perilaku evolusi skema yang fleksibel dan dapat dikonfigurasi.
-
Penyaringan data menggunakan klausa
WHERE. -
Dukungan untuk pemangkasan kolom.
Flink CDC vs. Flink DataStream
Tabel berikut mencantumkan keunggulan pekerjaan ingesti data Flink CDC dibandingkan pekerjaan DataStream.
|
Flink CDC Ingestion |
Flink DataStream |
|
Dapat diakses oleh pengguna dengan berbagai tingkat keahlian. |
Memerlukan keahlian dalam Java dan sistem terdistribusi. |
|
Menyembunyikan kompleksitas dasar dan menyederhanakan pengembangan. |
Memerlukan pengetahuan tentang framework Flink. |
|
Format YAML yang mudah dipelajari. |
Memerlukan pengetahuan tentang alat seperti Maven untuk manajemen dependensi. |
|
Daya guna ulang tinggi untuk pekerjaan yang sudah ada. |
Sulit menggunakan ulang kode yang sudah ada. |
Batasan
-
Gunakan Ververica Runtime (VVR) 11.1 atau versi yang lebih baru untuk mengembangkan pekerjaan ingesti data Flink CDC. Jika Anda perlu menggunakan VVR 8.x, gunakan VVR 8.0.11.
-
Setiap pekerjaan hanya mendukung satu sumber dan satu sink. Untuk membaca dari beberapa sumber atau menulis ke beberapa sink, Anda harus membuat beberapa pekerjaan Flink CDC.
-
Pekerjaan Flink CDC tidak dapat dideploy ke kluster session.
-
Pekerjaan ingesti data Flink CDC tidak mendukung tuning otomatis.
Konektor ingesti data Flink CDC
Untuk detail tentang konektor sumber dan sink yang didukung untuk ingesti data Flink CDC, lihat Konektor yang didukung.
Membuat pekerjaan ingesti data Flink CDC
Dari templat
-
Masuk ke Konsol Realtime Compute for Apache Flink.
-
Pada kolom Actions untuk ruang kerja target, klik Console.
-
Di panel navigasi sebelah kiri, pilih .
-
Klik
lalu klik New Draft with Template. -
Pilih templat sinkronisasi data.
Saat ini, hanya tersedia templat MySQL ke StarRocks, MySQL ke Paimon, dan MySQL ke Hologres.
-
Masukkan informasi pekerjaan, seperti OK, location, dan engine version, lalu klik OK.
-
Konfigurasikan informasi sumber dan sink untuk pekerjaan Flink CDC.
Untuk detail konfigurasi parameter, lihat dokumentasi konektor terkait.
Dari pekerjaan CTAS/CDAS
-
Jika suatu pekerjaan berisi beberapa pernyataan CXAS, Flink hanya mendeteksi dan mengonversi pernyataan pertama.
-
Karena perbedaan dukungan fungsi bawaan antara Flink SQL dan Flink CDC, aturan
transformyang dihasilkan mungkin tidak berfungsi sebagaimana mestinya. Anda harus meninjau dan menyesuaikannya sesuai kebutuhan. -
Jika sumbernya adalah MySQL dan pekerjaan CTAS/CDAS asli masih berjalan, Anda harus menyesuaikan
server-iduntuk sumber dalam pekerjaan ingesti data Flink CDC agar tidak terjadi konflik.
-
Masuk ke Konsol Realtime Compute for Apache Flink.
-
Pada kolom Actions untuk ruang kerja target, klik Console.
-
Di panel navigasi sebelah kiri, pilih .
-
Klik
, lalu klik New Draft from CTAS/CDAS. Pilih pekerjaan CTAS atau CDAS target, lalu klik OK.Di halaman pemilihan, sistem hanya menampilkan pekerjaan CTAS dan CDAS yang valid. Pekerjaan ETL biasa dan draft dengan kesalahan sintaks tidak ditampilkan.
-
Masukkan informasi pekerjaan, seperti OK, location, dan engine version, lalu klik OK.
Dari Flink CDC open-source
-
Masuk ke Konsol Realtime Compute for Apache Flink.
-
Pada kolom Actions untuk ruang kerja target, klik Console.
-
Di panel navigasi sebelah kiri, pilih .
-
Klik
, lalu pilih New Draft. Masukkan name dan engine version, lalu klik Create. -
Tempel kode YAML dari pekerjaan Flink CDC open-source Anda ke editor.
-
(Opsional) Klik Validate.
Opsi ini memeriksa kesalahan sintaks, masalah konektivitas jaringan, dan masalah izin.
Dari awal
-
Masuk ke Konsol Realtime Compute for Apache Flink.
-
Pada kolom Actions untuk ruang kerja target, klik Console.
-
Di panel navigasi sebelah kiri, pilih .
-
Klik
, lalu pilih New Draft. Masukkan name dan engine version, lalu klik Create. -
Konfigurasikan pekerjaan Flink CDC menggunakan YAML. Contoh:
# Wajib source: # Tipe sumber data. type: <Ganti dengan tipe konektor sumber Anda> # Konfigurasi untuk sumber data. Untuk detail item konfigurasi, lihat dokumentasi konektor terkait. ... # Wajib sink: # Tipe sink. type: <Ganti dengan tipe konektor sink Anda> # Konfigurasi untuk sink. Untuk detail item konfigurasi, lihat dokumentasi konektor terkait. ... # Opsional transform: # Aturan transformasi untuk tabel flink_test.customers. - source-table: flink_test.customers # Konfigurasi proyeksi. Menentukan kolom yang akan disinkronkan dan melakukan transformasi data. projection: id, username, UPPER(username) as username1, age, (age + 1) as age1, test_col1, __schema_name__ || '.' || __table_name__ identifier_name # Kondisi filter. Hanya menyinkronkan data dengan id lebih besar dari 10. filter: id > 10 # Deskripsi aturan transformasi. description: tambahkan kolom terhitung berdasarkan tabel sumber # Opsional route: # Aturan entri rute yang menentukan pemetaan antara tabel sumber dan tabel sink. - source-table: flink_test.customers sink-table: db.customers_o # Deskripsi aturan entri rute. description: sinkronisasi tabel customers - source-table: flink_test.customers_suffix sink-table: db.customers_s # Deskripsi aturan entri rute. description: sinkronisasi tabel customers_suffix # Opsional pipeline: # Nama pekerjaan. name: MySQL to Hologres PipelineCatatanDalam pekerjaan Flink CDC, pisahkan kunci dan nilainya dengan tanda titik dua dan spasi. Formatnya adalah
Key: Value.Tabel berikut menjelaskan blok kode tersebut.
Wajib
Modul
Deskripsi
Ya
source
Titik awal pipa data. Flink CDC menangkap data perubahan dari sumber.
Catatan-
Saat ini, MySQL adalah satu-satunya sumber yang didukung. Untuk detail konfigurasi spesifik, lihat Konektor MySQL.
-
Anda dapat menggunakan variabel untuk mengelola informasi sensitif. Untuk informasi lebih lanjut, lihat Manajemen Variabel.
sink
Titik akhir pipa data. Flink CDC mengirimkan perubahan data yang ditangkap ke sistem sink.
Catatan-
Untuk informasi tentang sink yang didukung, lihat Konektor ingesti data Flink CDC. Untuk detail konfigurasi, lihat dokumentasi konektor spesifik.
-
Anda dapat menggunakan variabel untuk mengelola informasi sensitif. Untuk informasi lebih lanjut, lihat Manajemen Variabel.
Tidak
pipeline
(pipa data)
Menentukan konfigurasi dasar untuk seluruh pekerjaan pipa data, seperti nama pipeline.
transform
Menentukan aturan transformasi data untuk mengoperasikan data saat mengalir melalui pipeline Flink. Mendukung pemrosesan ETL, penyaringan klausa
WHERE, pemangkasan kolom, dan kolom terhitung.Jika Anda perlu mentransformasi data perubahan mentah yang ditangkap oleh Flink CDC agar sesuai dengan sistem downstream tertentu, gunakan blok
transform.route
Jika modul ini tidak dikonfigurasi, pekerjaan secara default melakukan sinkronisasi seluruh database atau tabel target.
Dalam beberapa kasus, Anda mungkin perlu mengirim data perubahan yang ditangkap ke tujuan berbeda berdasarkan aturan tertentu. Modul
routememungkinkan Anda secara fleksibel menentukan hubungan pemetaan antara sumber dan sink, serta mengirim data ke target berbeda.Untuk informasi lebih lanjut tentang sintaks dan konfigurasi setiap modul, lihat Referensi pekerjaan ingesti data Flink CDC.
Kode berikut memberikan contoh sinkronisasi semua tabel dari database app_db di MySQL ke database di Hologres.
source: type: mysql hostname: <hostname> port: 3306 username: ${secret_values.mysqlusername} password: ${secret_values.mysqlpassword} tables: app_db.\.* server-id: 5400-5404 # (Opsional) Sinkronkan data dari tabel yang baru dibuat pada fase inkremental. scan.binlog.newly-added-table.enabled: true # (Opsional) Sinkronkan komentar tabel dan bidang. include-comments.enabled: true # (Opsional) Utamakan pengiriman chunk tak terbatas untuk mencegah potensi masalah OutOfMemory Pengelola Tugas. 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: <endpoint> dbname: <database-name> username: ${secret_values.holousername} password: ${secret_values.holopassword} pipeline: name: Sync MySQL Database to Hologres -
-
(Opsional) Klik Validate.
Opsi ini memeriksa kesalahan sintaks, masalah konektivitas jaringan, dan masalah izin.
Topik terkait
-
Setelah selesai mengembangkan pekerjaan Flink CDC, Anda perlu mendeploynya. Untuk informasi lebih lanjut, lihat Deploy a job.
-
Untuk membangun pekerjaan Flink CDC secara cepat guna menyinkronkan data dari database MySQL ke StarRocks, lihat Tutorial: Create a Flink CDC data ingestion job.