All Products
Search
Document Center

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

Last Updated:Apr 23, 2026

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 SchemaChangeEvent untuk mentransmisikan informasi evolusi skema, seperti pembuatan tabel, penambahan kolom, atau pemotongan tabel (truncate). Mereka menggunakan DataChangeEvent untuk 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 CREATE TABLE dan INSERT

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 UPDATE tidak dipisah.

  • Menyinkronkan lebih banyak jenis evolusi skema, seperti TRUNCATE TABLE dan DROP 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

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

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

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

  4. Klik image lalu klik New from Template.

  5. Pilih templat sinkronisasi data.

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

    image

  6. Masukkan informasi pekerjaan, termasuk nama pekerjaan, lokasi penyimpanan, dan versi engine, lalu klik OK.

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

    Lihat dokumentasi konektor terkait untuk detail konfigurasi parameter.

Dari CTAS/CDAS

Penting
  • Jika suatu pekerjaan berisi beberapa pernyataan CTAS atau CDAS, 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 langsung. Anda harus meninjau dan menyesuaikannya sesuai kebutuhan.

  • Jika sumbernya adalah MySQL dan pekerjaan CTAS/CDAS asli masih berjalan, Anda harus menyesuaikan server-id sumber pada pekerjaan ingesti data Flink CDC untuk menghindari konflik.

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

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

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

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

  5. Masukkan informasi pekerjaan, termasuk nama pekerjaan, lokasi penyimpanan, dan versi engine, lalu klik OK.

Dari open source

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

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

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

  4. Klik image, pilih New Data Ingestion Draft, masukkan File Name dan Engine Version, lalu klik Create.

  5. Salin kode pekerjaan Flink CDC open source.

  6. (Opsional) Klik Validate.

    Anda dapat memvalidasi sintaksis, konektivitas jaringan, dan izin akses.

Dari awal

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

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

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

  4. Klik image, pilih New Data Ingestion Draft, masukkan File Name dan Engine Version, lalu klik Create.

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

    Dalam 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.

    Catatan
    • Saat 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.

    Catatan
    • Untuk 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 transform untuk 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
  6. (Opsional) Klik Validate.

    Anda dapat memvalidasi sintaksis, konektivitas jaringan, dan izin akses.

Dokumen terkait