All Products
Search
Document Center

Realtime Compute for Apache Flink:Ingesti database waktu nyata

Last Updated:Apr 29, 2026

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

Siapkan data uji MySQL dan database Hologres

  1. Klik tpc_ds.sql, user_db1.sql, user_db2.sql, dan user_db3.sql untuk mengunduh file data uji ke mesin lokal Anda.

  2. Di Konsol DMS, siapkan data uji di instans ApsaraDB RDS for MySQL Anda.

    1. Masuk ke instans ApsaraDB RDS for MySQL Anda menggunakan DMS.

    2. Di jendela SQL Console, masukkan perintah berikut lalu klik Execute.

      Perintah berikut membuat empat database: tpc_ds, user_db1, user_db2, dan user_db3.

      CREATE DATABASE tpc_ds;
      CREATE DATABASE user_db1;
      CREATE DATABASE user_db2;
      CREATE DATABASE user_db3;
    3. Di bilah navigasi atas, klik Data Import.

    4. 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, dan user_db3.导入数据

  3. Di Konsol Hologres, buat database bernama my_user untuk 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.

  1. Dapatkan Blok CIDR ruang kerja Flink.

    1. Masuk ke Konsol Realtime Compute for Apache Flink.

    2. Di daftar ruang kerja, temukan workspace target lalu pilih More > Workspace Details di kolom Actions.

    3. Di dialog Workspace Details, lihat informasi CIDR Block untuk vSwitch Flink.

      网段信息

  2. Tambahkan Blok CIDR Flink ke daftar putih IP instans ApsaraDB RDS for MySQL.

    Untuk informasi lebih lanjut, lihat Konfigurasikan daftar putih IP.RDS白名单

  3. 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.Holo白名单

Langkah 1: Kembangkan pekerjaan ingesti

  1. Masuk ke Konsol pengembangan Flink dan buat pekerjaan baru.

    1. Di halaman Data Development > Data Ingestion, klik New.

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

    3. Klik Next.

    4. Di dialog New Data Ingestion Job Draft, konfigurasikan parameter pekerjaan.

      Parameter

      Deskripsi

      Contoh

      File Name

      Nama pekerjaan.

      Catatan

      Nama 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

    5. Klik OK.

  2. Salin kode pekerjaan berikut ke editor pekerjaan.

    Kode berikut menyinkronkan semua tabel dari database tpc_ds ke 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.users
    Catatan

    Tabel dari database MySQL tpc_ds dipetakan langsung ke tabel dengan nama yang sama di tujuan, sehingga bagian route tidak memerlukan konfigurasi pemetaan tambahan. Untuk menyinkronkan tabel ke database dengan nama berbeda, seperti ods_tps_ds, konfigurasikan modul route sebagai 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

  1. Di halaman Data Development > Data Ingestion, klik Deploy. Di dialog yang muncul, klik Confirm.部署

  2. Di halaman Operation Center > Job Operations, klik Start di kolom Actions untuk pekerjaan target. Konfigurasikan parameter sesuai kebutuhan. Untuk informasi lebih lanjut, lihat Mulai pekerjaan.

  3. Klik Start.

    Setelah pekerjaan dimulai, Anda dapat melihat informasi waktu proses dan statusnya di halaman Job Operations.作业状态

Langkah 3: Verifikasi sinkronisasi penuh

  1. Masuk ke Konsol Manajemen Hologres.

  2. Di tab Metadata Management, verifikasi bahwa 24 tabel beserta datanya tersedia di database tpc_ds instans Hologres.

    holo表数据

  3. Di tab Metadata Management, periksa skema tabel users di database my_user.

    Gambar berikut menunjukkan skema dan data tabel yang telah disinkronkan.

    • Skema tabel表结构

      Skema tabel users mencakup dua kolom tambahan yang tidak ada di tabel MySQL sumber: _db_name dan _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.

  1. Masuk ke Konsol Realtime Compute for Apache Flink.

  2. Klik Console di kolom Actions untuk ruang kerja target.

  3. Di halaman Operation Center > Job Operations, klik nama pekerjaan target.

  4. Klik tab Monitoring and Alerts (atau Metrics).

  5. Periksa grafik currentEmitEventTimeLag untuk menentukan fase sinkronisasi data.

    数据曲线

    • Nilai 0 menunjukkan fase sinkronisasi penuh.

    • Nilai lebih dari 0 menunjukkan fase sinkronisasi inkremental.

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

    1. Masuk ke instans ApsaraDB RDS for MySQL menggunakan DMS.

    2. Di database user_db2, jalankan perintah berikut untuk mengubah skema tabel user02 serta 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.
    3. 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 user02 dan modifikasi data disebarkan secara waktu nyata, meskipun tabel terbagi memiliki skema berbeda. Tabel users Hologres kini menampilkan kolom age baru, 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.

  1. Di halaman Operation Center > Job Operations, klik nama pekerjaan target.

  2. Di tab Deployment Details, klik Edit di pojok kanan atas bagian Resource Configuration.

  3. Atur secara manual parameter sumber daya seperti memori TaskManager dan konkurensi.

  4. Di sisi kanan bagian Resource Configuration, klik Save.

  5. Restart pekerjaan.

    Perubahan konfigurasi sumber daya hanya berlaku setelah Anda me-restart pekerjaan.

Dokumen terkait