All Products
Search
Document Center

Realtime Compute for Apache Flink:Mengembangkan pekerjaan ingesti data Flink CDC

Last Updated:Aug 07, 2026

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 objek RowData memiliki tipe perubahan sendiri, yang mencakup empat tipe utama: insert (+I), update before (-U), update after (+U), dan delete (-D).

  • Flink CDC menggunakan SchemaChangeEvent untuk menyampaikan informasi perubahan skema, seperti membuat tabel, menambah kolom, atau memotong (truncate) tabel. Flink CDC menggunakan DataChangeEvent untuk 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 CREATE TABLE dan INSERT

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 TABLE dan DROP 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

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

  2. Pada kolom Actions untuk ruang kerja target, klik Console.

  3. Di panel navigasi sebelah kiri, pilih Development > Data Ingestion.

  4. Klik image lalu klik New Draft with Template.

  5. Pilih templat sinkronisasi data.

    Saat ini, hanya tersedia templat MySQL ke StarRocks, MySQL ke Paimon, dan MySQL ke Hologres.

  6. Masukkan informasi pekerjaan, seperti OK, location, dan engine version, lalu klik OK.

  7. Konfigurasikan informasi sumber dan sink untuk pekerjaan Flink CDC.

    Untuk detail konfigurasi parameter, lihat dokumentasi konektor terkait.

Dari pekerjaan CTAS/CDAS

Penting
  • 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 transform yang 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-id untuk sumber dalam pekerjaan ingesti data Flink CDC agar tidak terjadi konflik.

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

  2. Pada kolom Actions untuk ruang kerja target, klik Console.

  3. Di panel navigasi sebelah kiri, pilih Development > Data Ingestion.

  4. Klik image, 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.

  5. Masukkan informasi pekerjaan, seperti OK, location, dan engine version, lalu klik OK.

Dari Flink CDC open-source

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

  2. Pada kolom Actions untuk ruang kerja target, klik Console.

  3. Di panel navigasi sebelah kiri, pilih Development > Data Ingestion.

  4. Klik image, lalu pilih New Draft. Masukkan name dan engine version, lalu klik Create.

  5. Tempel kode YAML dari pekerjaan Flink CDC open-source Anda ke editor.

  6. (Opsional) Klik Validate.

    Opsi ini memeriksa kesalahan sintaks, masalah konektivitas jaringan, dan masalah izin.

Dari awal

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

  2. Pada kolom Actions untuk ruang kerja target, klik Console.

  3. Di panel navigasi sebelah kiri, pilih Development > Data Ingestion.

  4. Klik image, lalu pilih New Draft. Masukkan name dan engine version, lalu klik Create.

  5. 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 Pipeline
    Catatan

    Dalam 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 route memungkinkan 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
  6. (Opsional) Klik Validate.

    Opsi ini memeriksa kesalahan sintaks, masalah konektivitas jaringan, dan masalah izin.

Topik terkait