All Products
Search
Document Center

Realtime Compute for Apache Flink:Ingesti Database Real-time

Last Updated:Aug 07, 2026

Realtime Compute for Apache Flink menyederhanakan ingesti data real-time 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

Berikut ini menunjukkan struktur database dan tabel di konsol DMS.

Instans RDS berisi beberapa database (tpc_ds, tpc_ds_large, user_db1, user_db2, user_db3, __recycle_bin__). Database user_db3 berisi tiga tabel: user03, user06, dan user09. Tabel user03 memiliki dua kolom: id (int(11)) dan name (varchar(255)).

Ikuti langkah-langkah berikut untuk mengembangkan pekerjaan ingesti data yang menyinkronkan semua tabel tersebut ke Hologres dan menggabungkan tabel user terpartisi menjadi satu tabel:

Topik ini menggunakan ingesti data Flink CDC untuk melakukan sinkronisasi seluruh database dan menggabungkan tabel terpartisi. Pendekatan ini memungkinkan Anda menyelesaikan sinkronisasi data penuh dan inkremental, serta sinkronisasi perubahan skema real-time, 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 dan 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 operasi ini untuk mengimpor file data yang sesuai ke database tpc_ds, user_db1, user_db2, dan user_db3.

      Atur File Encoding ke Auto Detect, Import Mode ke Express Mode atau Safe Mode, dan File Type ke SQL Script, CSV, atau Excel. Lampiran mendukung format txt/sql/csv/xlsx/zip, hingga 5 GB.

  3. Di konsol Hologres, buat database bernama my_user untuk menyimpan data tabel user 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 dan pilih More > Workspace Details di kolom Actions.

    3. Di kotak 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.

    Di kotak dialog Modify Whitelist Group, masukkan Blok CIDR Flink fully managed di bidang Group Whitelist, lalu klik OK.

  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.

    Di HoloWeb, pilih Security Center dari bilah navigasi atas, klik IP Whitelist di panel kiri, dan konfigurasikan parameter berikut di kotak dialog Edit IP Whitelist:

    • Group: Pilih default

    • Database Restriction: Pilih ALL

    • User Restriction: Pilih ALL

    • IP Address: Masukkan Blok CIDR Flink dalam Notasi CIDR (misalnya, 172.xx.0/19). Pisahkan beberapa alamat IP dengan baris baru.

    Klik OK untuk menerapkan.

Langkah 1: Kembangkan pekerjaan ingesti

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

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

    2. Klik Blank Data Ingestion Draft.

      Realtime Compute for Apache Flink menyediakan berbagai templat kode, masing-masing 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 kotak dialog New Draft, konfigurasikan parameter.

      Parameter

      Deskripsi

      Contoh

      Name

      Nama pekerjaan.

      Catatan

      Nama pekerjaan harus unik dalam proyek saat ini.

      flink-test

      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 engine, kompatibilitas versi, dan tanggal penting siklus hidup, lihat Engine Versions.

      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 user terpartisi 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
      # (Optional) Synchronize table and column comments.
      include-comments.enabled: true
      # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory errors.
      scan.incremental.snapshot.unbounded-chunk-first.enabled: true
      # (Optional) Enable parsing filters to accelerate reading.
      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:
      # Merge and synchronize the sharded user tables to the my_user.users table.
      - 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:
      # Merge and synchronize the sharded user tables to the my_user.users table.
      - source-table: user_db[0-9]+.user[0-9]+
        sink-table: my_user.users
      # Rename the database for all tables under tpc_ds and synchronize them to ods_tps_ds.
      - source-table: tpc_ds.\.*
        sink-table: ods_tps_ds.<>
        replace-symbol: <>

Langkah 2: Mulai pekerjaan

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

    Di kotak dialog, Anda dapat memasukkan Remarks, mengatur Job Tags, dan memilih antrian dari daftar drop-down Deployment Target (default: default-queue). Catatan: Penerapan berlaku saat pekerjaan dimulai berikutnya.

  2. Di halaman O&M > Deployments, 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. Status pekerjaan mencakup Running, Failed, dan Stopped. Gunakan filter drop-down Stream Jobs di bagian atas untuk memfilter daftar pekerjaan berdasarkan jenis.

Langkah 3: Verifikasi sinkronisasi penuh

  1. Masuk ke Hologres Management Console.

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

    Navigasi ke Instance > tpc_ds > public > Tables untuk melihat tabel seperti call_center, catalog_sales, customer, dan store_sales. Pilih tabel store_sales dan klik tab Data Preview untuk melihat kolom seperti ss_sold_date, ss_sold_time, ss_item_sk, dan ss_customer.

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

    Berikut ini menjelaskan skema dan data tabel yang disinkronkan.

    • Skema tabel

      Setelah sinkronisasi, tabel users berisi kolom-kolom berikut:

      • _db_name (text, primary key)

      • _table_name (text, primary key)

      • id (int4, primary key)

      • name (varchar 255, nullable)

      Kolom _db_name, _table_name, dan id membentuk kunci utama gabungan.

      Skema tabel users mencakup dua kolom tambahan yang tidak ada di tabel sumber MySQL: _db_name dan _table_name. Kolom-kolom ini menunjukkan database dan tabel sumber untuk setiap baris dan merupakan bagian dari kunci utama gabungan, memastikan keunikan data setelah tabel terpartisi 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 menunjukkan bahwa data user telah sepenuhnya disinkronkan ke tiga database: user_db1 (user01, user04, user07), user_db2 (user02, user05, user08), dan user_db3 (user03, user06, user09), total 9 catatan dengan kolom _db_name, _table_name, id, dan name.

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 O&M > Deployments, klik nama pekerjaan target.

  4. Klik tab Alarm (atau Metrics).

  5. Periksa grafik currentEmitEventTimeLag untuk menentukan fase sinkronisasi data.

    数据曲线

    • Nilai 0 menunjukkan fase sinkronisasi penuh.

    • Nilai lebih besar dari 0 menunjukkan fase sinkronisasi inkremental.

  6. Verifikasi sinkronisasi data dan perubahan skema real-time.

    Sumber CDC MySQL mendukung sinkronisasi data dan skema real-time selama fase inkremental. Untuk memverifikasi hal ini, ubah skema dan data tabel user terpartisi 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;   -- Add the age column.
      INSERT INTO `user02` (id, name, age) VALUES (27, 'Tony', 30); -- Insert a row that includes the age data.
      UPDATE `user05` SET name='JARK' WHERE id=15;  -- Update another table and change the name to uppercase.
    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;

      Perubahan skema pada user02 dan modifikasi data disebarkan secara real-time, meskipun tabel terpartisi memiliki skema berbeda. Tabel users Hologres kini menampilkan kolom age baru, catatan yang disisipkan untuk Tony, dan catatan yang diperbarui untuk JARK.

      Hasil kueri menunjukkan 11 baris dari tabel users di user_db1, user_db2, dan user_db3. Perubahan sinkronisasi inkremental tercermin: user Tony (id 27) di user_db2 memiliki age diperbarui menjadi 30, dan catatan lain memiliki name diubah menjadi JARK. Baris-baris lain memiliki nilai age \N, yang mengonfirmasi bahwa sinkronisasi inkremental berfungsi.

(Opsional) Langkah 5: Konfigurasikan resource pekerjaan

Untuk kinerja yang lebih baik, Anda dapat menyesuaikan resource pekerjaan seperti konkurensi, memori TaskManager, dan CUs berdasarkan volume data Anda.

  1. Di halaman O&M > Deployments, klik nama pekerjaan target.

  2. Di tab Configuration, klik Edit di pojok kanan atas bagian Resources.

  3. Atur manual parameter resource seperti memori TaskManager dan konkurensi.

  4. Di sisi kanan bagian Resources, klik Save.

  5. Restart pekerjaan.

    Perubahan konfigurasi resource hanya berlaku setelah Anda me-restart pekerjaan.

Dokumen terkait