Realtime Compute for Apache Flink menggunakan Flink CDC untuk melakukan ingesti data. Anda dapat mengembangkan pekerjaan YAML guna menyinkronkan data secara efisien dari sumber ke sink. Topik ini menjelaskan cara mengembangkan pekerjaan ingesti data Flink CDC.
Informasi latar belakang
Ingesti data Flink CDC menggunakan Flink CDC untuk menyederhanakan integrasi data. Dengan menggunakan YAML untuk mendefinisikan proses ETL kompleks yang secara otomatis dikonversi menjadi logika runtime Flink, Anda dapat menerapkan fitur seperti sinkronisasi database penuh, sinkronisasi tabel terpartisi, evolusi skema, dan kolom terhitung. Pendekatan ini sangat 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 pekerjaan ingesti data Flink CDC dibandingkan dua metode pengembangan lainnya.
Flink CDC vs. Flink SQL
Pekerjaan ingesti data Flink CDC dan pekerjaan Flink SQL mentransmisikan jenis data yang berbeda:
Pekerjaan SQL mentransmisikan
RowData, di mana setiap baris memiliki tipe perubahan: insert (+I), update before (-U), update after (+U), atau delete (-D).Pekerjaan Flink CDC menggunakan
SchemaChangeEventuntuk mentransmisikan informasi evolusi skema, seperti pembuatan tabel, penambahan kolom, atau pemotongan tabel (truncate). Mereka menggunakanDataChangeEventuntuk mentransmisikan perubahan data seperti insert, update, dan delete. Pesan update berisi status sebelum dan sesudah suatu baris, sehingga memungkinkan Anda menulis data perubahan asli ke sink.
Tabel berikut mencantumkan keunggulan pekerjaan ingesti data Flink CDC dibandingkan pekerjaan SQL.
Flink CDC | Flink SQL |
Secara otomatis menemukan skema dan mendukung sinkronisasi database penuh | Memerlukan pernyataan manual |
Mendukung berbagai kebijakan evolusi skema | Tidak mendukung evolusi skema |
Mempertahankan changelog asli | Mengganggu struktur changelog asli |
Membaca dan menulis ke beberapa tabel | Membaca dan menulis ke satu tabel saja |
Dibandingkan dengan pernyataan CTAS atau CDAS, pekerjaan Flink CDC menawarkan keunggulan berikut:
Langsung menyinkronkan evolusi skema hulu tanpa menunggu penulisan data baru untuk memicu proses tersebut.
Mempertahankan changelog asli, dan pesan
UPDATEtidak dipisah.Menyinkronkan lebih banyak jenis evolusi skema, seperti
TRUNCATE TABLEdanDROP TABLE.Mendukung pemetaan tabel untuk secara fleksibel menentukan nama tabel sink.
Mendukung perilaku evolusi skema yang fleksibel dan dapat dikonfigurasi pengguna.
Mendukung penyaringan data dengan menggunakan klausa
WHERE.Mendukung pemangkasan kolom (column pruning).
Flink CDC vs. Flink DataStream
Tabel berikut mencantumkan keunggulan pekerjaan ingesti data Flink CDC dibandingkan pekerjaan DataStream.
Flink CDC | Flink DataStream |
Dirancang untuk pengguna dengan berbagai tingkat keahlian, bukan hanya ahli | Memerlukan keahlian dalam Java dan sistem terdistribusi |
Menyembunyikan detail tingkat rendah dan menyederhanakan pengembangan | Memerlukan keakraban dengan framework Flink |
Format YAML mudah dipahami dan dipelajari | Memerlukan alat seperti Maven untuk mengelola dependensi |
Pekerjaan yang sudah ada mudah digunakan kembali | Kode yang sudah ada sulit digunakan kembali |
Batasan
Gunakan Ververica Runtime (VVR) 11.1 atau versi yang lebih baru untuk mengembangkan pekerjaan ingesti data Flink CDC. Jika Anda perlu menggunakan versi 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 diterapkan ke session cluster.
Pekerjaan Flink CDC tidak mendukung penyetelan otomatis.
Flink CDC konektor ingesti data
Lihat Konektor yang didukung untuk daftar konektor yang didukung sebagai sumber dan sink untuk ingesti data Flink CDC.
Membuat pekerjaan ingesti data Flink CDC
Dari templat
Masuk ke Konsol Realtime Compute for Apache Flink.
Pada kolom Actions ruang kerja target, klik Console.
Di panel navigasi kiri, pilih .
Klik
lalu klik New from Template.Pilih templat sinkronisasi data.
Saat ini tersedia templat untuk MySQL ke StarRocks, MySQL ke Paimon, dan MySQL ke Hologres.

Masukkan informasi pekerjaan, termasuk nama pekerjaan, lokasi penyimpanan, dan versi engine, lalu klik OK.
Konfigurasikan informasi sumber dan sink untuk pekerjaan Flink CDC.
Lihat dokumentasi konektor terkait untuk detail konfigurasi parameter.
Dari CTAS/CDAS
Jika suatu pekerjaan berisi beberapa pernyataan
CTASatauCDAS, Flink hanya mendeteksi dan mengonversi pernyataan pertama.Karena perbedaan dukungan fungsi bawaan antara Flink SQL dan Flink CDC, aturan
transformyang dihasilkan mungkin tidak berfungsi langsung. Anda harus meninjau dan menyesuaikannya sesuai kebutuhan.Jika sumbernya adalah MySQL dan pekerjaan
CTAS/CDASasli masih berjalan, Anda harus menyesuaikanserver-idsumber pada pekerjaan ingesti data Flink CDC untuk menghindari konflik.
Masuk ke Konsol Realtime Compute for Apache Flink.
Pada kolom Actions ruang kerja target, klik Console.
Di panel navigasi kiri, pilih .
Klik
lalu klik New from CTAS/CDAS Job. Pilih pekerjaan CTAS atau CDAS target dan klik OK.Pada halaman pemilihan, hanya pekerjaan CTAS dan CDAS yang valid yang ditampilkan. Halaman ini tidak menampilkan pekerjaan ETL biasa atau draf dengan kesalahan sintaksis.
Masukkan informasi pekerjaan, termasuk nama pekerjaan, lokasi penyimpanan, dan versi engine, lalu klik OK.
Dari open source
Masuk ke Konsol Realtime Compute for Apache Flink.
Pada kolom Actions ruang kerja target, klik Console.
Di panel navigasi kiri, pilih .
Klik
, pilih New Data Ingestion Draft, masukkan File Name dan Engine Version, lalu klik Create.Salin kode pekerjaan Flink CDC open source.
(Opsional) Klik Validate.
Anda dapat memvalidasi sintaksis, konektivitas jaringan, dan izin akses.
Dari awal
Masuk ke Konsol Realtime Compute for Apache Flink.
Pada kolom Actions ruang kerja target, klik Console.
Di panel navigasi kiri, pilih .
Klik
, pilih New Data Ingestion Draft, masukkan File Name dan Engine Version, lalu klik Create.Konfigurasikan pekerjaan Flink CDC.
# Wajib source: # Jenis konektor sumber type: <Ganti dengan jenis konektor sumber Anda> # Konfigurasi sumber. Untuk detailnya, lihat dokumentasi konektor terkait. ... # Wajib sink: # Jenis konektor sink type: <Ganti dengan jenis konektor sink Anda> # Konfigurasi sink. Untuk detailnya, lihat dokumentasi konektor terkait. ... # Opsional transform: # Aturan transformasi untuk tabel flink_test.customers - source-table: flink_test.customers # Pengaturan 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 routing. Menentukan pemetaan antara tabel sumber dan tabel sink. - source-table: flink_test.customers sink-table: db.customers_o # Deskripsi aturan routing description: sinkronisasi tabel customers - source-table: flink_test.customers_suffix sink-table: db.customers_s # Deskripsi aturan routing description: sinkronisasi tabel customers_suffix # Opsional pipeline: # Nama pekerjaan name: MySQL to Hologres PipelineCatatanDalam pekerjaan Flink CDC, key dan value harus dipisahkan oleh spasi dalam format
Key: Value.Tabel berikut menjelaskan blok kode tersebut.
Wajib
Modul
Deskripsi
Wajib
source
Awal pipa data. Flink CDC menangkap data perubahan dari sumber.
CatatanSaat ini, MySQL adalah satu-satunya sumber yang didukung. Untuk detail item konfigurasi tertentu, lihat MySQL.
Anda dapat menggunakan variabel untuk mengelola informasi sensitif. Untuk informasi selengkapnya, lihat Manajemen Variabel.
sink
Akhir pipa data. Flink CDC mentransmisikan perubahan data yang ditangkap ke sink.
CatatanUntuk sistem sink yang didukung, lihat Konektor ingesti data Flink CDC. Untuk detail item konfigurasi sink, lihat dokumentasi konektor terkait.
Anda dapat menggunakan variabel untuk mengelola informasi sensitif. Untuk informasi selengkapnya, lihat Manajemen Variabel.
Opsional
pipeline
(pipa data)
Menentukan konfigurasi dasar untuk seluruh pekerjaan pipa data, seperti nama pipeline.
transform (transformasi data)
Menentukan aturan transformasi data. Transformasi beroperasi pada data saat mengalir melalui pipeline Flink. Fitur yang didukung mencakup pemrosesan ETL, penyaringan dengan klausa
WHERE, pemangkasan kolom, dan kolom terhitung.Gunakan modul
transformuntuk mengubah data perubahan mentah dari Flink CDC agar sesuai dengan sistem downstream tertentu.route
Jika modul ini tidak dikonfigurasi, pekerjaan akan melakukan sinkronisasi database penuh atau tabel tertentu.
Dalam beberapa kasus, Anda mungkin perlu mengirim data perubahan yang ditangkap ke tujuan berbeda berdasarkan aturan tertentu. Mekanisme routing memungkinkan Anda secara fleksibel menentukan hubungan pemetaan antara sumber dan sink untuk mengirim data ke sink yang berbeda.
Untuk detail sintaksis dan konfigurasi setiap modul, lihat Referensi untuk mengembangkan pekerjaan ingesti data Flink CDC.
Kode berikut memberikan contoh cara menyinkronkan 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 dibuat selama fase inkremental. scan.binlog.newly-added-table.enabled: true # (Opsional) Sinkronkan komentar tabel dan kolom. include-comments.enabled: true # (Opsional) Utamakan pengiriman chunk tak terbatas untuk menghindari potensi masalah OutOfMemory pada TaskManager. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Opsional) Aktifkan penyaringan parse 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.
Anda dapat memvalidasi sintaksis, konektivitas jaringan, dan izin akses.
Dokumen terkait
Setelah mengembangkan pekerjaan Flink CDC, Anda harus menerapkannya. Lihat Menerapkan pekerjaan untuk instruksi penerapan.
Untuk membangun pekerjaan Flink CDC yang menyinkronkan data dari database MySQL ke StarRocks secara cepat, lihat Panduan cepat untuk pekerjaan ingesti data Flink CDC.