All Products
Search
Document Center

AnalyticDB:Berlangganan log biner menggunakan Flink

Last Updated:Jul 10, 2026

Realtime Compute for Apache Flink berlangganan ke AnalyticDB for MySQL dari kluster AnalyticDB for MySQL untuk menangkap dan memproses perubahan database secara real time, sehingga memungkinkan sinkronisasi data dan komputasi aliran yang efisien.

Prasyarat

  • Kluster AnalyticDB for MySQL Anda harus merupakan salah satu edisi berikut: Enterprise Edition, Basic Edition, Data Lakehouse Edition, atau Data Warehouse Edition (dalam mode elastis).

  • Kluster AnalyticDB for MySQL harus memenuhi persyaratan versi kernel berikut:

    • xuanwu_v1 table engine: 3.2.1.0 atau yang lebih baru.

    • xuanwu_v2 table engine: 3.2.6.0 atau yang lebih baru.

    • Untuk berlangganan log biner dari tampilan yang di-materialisasi inkremental, versi kernel harus 3.2.6.9, 3.2.7.1, atau yang lebih baru.

    Catatan

    Untuk melihat dan memperbarui versi minor, buka bagian Configuration Information pada halaman Cluster Information di Konsol AnalyticDB for MySQL.

  • Ruang kerja Flink harus menggunakan Ververica Runtime (VVR) 8.0.4 atau versi yang lebih baru.

  • AnalyticDB for MySQL kluster dan ruang kerja fully managed Flink harus berada dalam VPC yang sama.

  • Tambahkan Blok CIDR dari ruang kerja Flink ke daftar putih AnalyticDB for MySQL.

Batasan

  • Flink hanya dapat memproses tipe data dasar dan tipe data kompleks JSON dari log biner AnalyticDB for MySQL.

  • Operasi berikut tidak ditangkap saat Anda menggunakan Flink untuk berlangganan log biner AnalyticDB for MySQL. Perubahan data terkait tidak disinkronkan ke konsumen downstream:

Langkah 1: Aktifkan binary logging

  1. Aktifkan binary logging. Contoh ini menggunakan tabel bernama source_table.

    Catatan

    AnalyticDB for MySQL hanya mendukung pengaktifan binary logging pada tingkat tabel.

    Saat membuat tabel

    CREATE TABLE source_table (
      `id` INT,
      `num` BIGINT,
      PRIMARY KEY (`id`)
    )DISTRIBUTED BY HASH (id) BINLOG=true;

    Untuk tabel yang sudah ada

    ALTER TABLE source_table BINLOG=true;
  2. (Opsional) Ubah periode retensi log biner.

    Anda dapat mengubah parameter binlog_ttl untuk menyesuaikan periode retensi log biner. Nilai default parameter ini adalah 6 jam. Sebagai contoh, Anda dapat mengatur periode retensi log biner untuk tabel source_table menjadi 1 hari.

    ALTER TABLE source_table binlog_ttl='1d';

    Parameter binlog_ttl mendukung format berikut:

    • Milidetik: Nilai numerik. Contoh: 60 merepresentasikan 60 milidetik.

    • Detik: Angka + s. Contoh: 30s merepresentasikan 30 detik.

    • Jam: Angka diikuti h. Contoh: 2h merepresentasikan 2 jam.

    • Hari: Angka diikuti d. Contoh: 1d merepresentasikan 1 hari.

    Catatan
    • Periode retensi log biner maksimum adalah 365 hari untuk kluster dengan versi kernel berikut atau versi yang lebih baru dalam seri masing-masing: 3.2.1.9, 3.2.2.14, 3.2.3.8, 3.2.4.4, dan 3.2.5.1. Untuk kluster dengan versi kernel sebelumnya, periode retensi maksimum adalah 21 hari.

    • Kami menyarankan agar Anda mengatur periode retensi log biner ke nilai yang tidak kurang dari nilai default parameter binlog_ttl. Jika periode retensi terlalu singkat, file dapat di-purge, yang dapat memengaruhi sinkronisasi data.

    • Jika Anda perlu melihat periode retensi log biner saat ini, jalankan SHOW CREATE TABLE source_table;.

Langkah 2: Unggah konektor AnalyticDB for MySQL

  1. Unduh file JAR konektor: flink-sql-connector-adb-mysql-cdc-2.4-20260420.jar.

  2. Masuk ke Konsol Realtime Compute for Apache Flink.

  3. Pada tab Flink, temukan ruang kerja target dan klik Console di kolom Actions.

  4. Di panel navigasi kiri, klik Connectors.

  5. Pada halaman Connectors, klik Create Custom Connector.

  6. Unggah konektor yang telah diunduh dan klik Next.

  7. Klik Finish. Konektor kustom yang dibuat akan muncul dalam daftar konektor.

Langkah 3: Berlangganan log biner

  1. Masuk ke Konsol Realtime Compute for Apache Flink dan buat pekerjaan SQL.

  2. Buat tabel sumber yang terhubung ke AnalyticDB for MySQL dan membaca data log biner dari tabel tertentu (source_table).

    Catatan
    • Kunci primer yang didefinisikan dalam DDL Flink harus sesuai dengan kunci primer pada tabel fisik di kluster AnalyticDB for MySQL. Ini mencakup kolom kunci dan nama kunci primer. Jika tidak, data mungkin tidak akurat.

    • Tipe data Flink harus kompatibel dengan tipe data di AnalyticDB for MySQL. Untuk detail pemetaan tipe data, lihat Pemetaan tipe.

    CREATE TEMPORARY TABLE adb_source (
      `id` INT,
      `num` BIGINT,
      PRIMARY KEY (`id`) NOT ENFORCED
    ) WITH (
      'connector' = 'adb-mysql-cdc',
      'hostname' = 'amv-2zepb9n1l58ct01z50000****.ads.aliyuncs.com',
      'username' = 'testUser',
      'password' = 'Test12****',
      'database-name' = 'binlog',
      'table-name' = 'source_table'
    );

    Tabel berikut menjelaskan parameter WITH.

    Parameter

    Wajib

    Default

    Tipe

    Deskripsi

    connector

    Ya

    Tidak ada

    STRING

    Konektor yang digunakan.

    Ini adalah konektor kustom. Atur parameter ini ke adb-mysql-cdc.

    hostname

    Ya

    Tidak ada

    STRING

    Titik akhir AnalyticDB for MySQL dari kluster AnalyticDB for MySQL.

    username

    Ya

    Tidak ada

    STRING

    Akun database untuk kluster AnalyticDB for MySQL.

    password

    Ya

    Tidak ada

    STRING

    Kata sandi untuk akun database AnalyticDB for MySQL.

    database-name

    Ya

    Tidak ada

    STRING

    Nama database AnalyticDB for MySQL.

    Karena AnalyticDB for MySQL menerapkan binary logging tingkat tabel, Anda hanya dapat menentukan satu database.

    table-name

    Ya

    Tidak ada

    STRING

    Nama tabel dalam database AnalyticDB for MySQL.

    Karena AnalyticDB for MySQL menerapkan binary logging tingkat tabel, Anda hanya dapat menentukan satu tabel.

    port

    Tidak

    3306

    INTEGER

    Nomor port.

    scan.incremental.snapshot.enabled

    Tidak

    true

    BOOLEAN

    Menentukan apakah mekanisme pembacaan snapshot inkremental diaktifkan.

    Fitur ini diaktifkan secara default. Snapshot inkremental adalah mekanisme baru untuk membaca snapshot tabel. Dibandingkan dengan mekanisme snapshot sebelumnya, snapshot inkremental memiliki keuntungan berikut:

    • Sumber mendukung pembacaan konkuren selama pembacaan snapshot.

    • Sumber mendukung checkpoint tingkat chunk selama pembacaan snapshot.

    • Sumber tidak perlu mendapatkan izin kunci database sebelum membaca snapshot.

    scan.incremental.snapshot.chunk.size

    Tidak

    8096

    INTEGER

    Jumlah baris per chunk untuk snapshot tabel.

    scan.snapshot.fetch.size

    Tidak

    1024

    INTEGER

    Jumlah maksimum baris yang diambil sekaligus saat membaca snapshot tabel.

    scan.startup.mode

    Tidak

    initial

    STRING

    Mode startup untuk konsumsi data.

    Nilai yang valid:

    • initial (default): Melakukan snapshot awal tabel lalu membaca log biner terbaru.

    • earliest-offset: Melewati fase snapshot dan mulai membaca dari log biner paling awal yang tersedia.

    • specific-offset: Melewati fase snapshot dan mulai dari posisi log biner tertentu. Tentukan nama file dan offset log biner dengan mengatur parameter scan.startup.specific-offset.file dan scan.startup.specific-offset.pos.

    • latest-offset: Melewati fase snapshot dan hanya membaca perubahan yang terjadi setelah konektor dimulai.

    • timestamp: Melewati fase snapshot dan mulai membaca dari timestamp tertentu. Atur timestamp dalam milidetik (ms) menggunakan parameter scan.startup.timestamp-millis.

    Penting

    Jika Anda menggunakan mode startup earliest-offset, specific-offset, atau timestamp, pastikan skema tabel tidak berubah antara posisi awal yang ditentukan dan saat pekerjaan dimulai. Jika tidak, pekerjaan mungkin gagal karena perubahan skema.

    scan.startup.specific-offset.file

    Tidak

    Tidak ada

    STRING

    Dalam mode startup specific-offset, nama file log biner pada posisi startup.

    Untuk mendapatkan nama file log biner terbaru, jalankan pernyataan SHOW MASTER STATUS for table_name;.

    scan.startup.specific-offset.pos

    Tidak

    Tidak ada

    LONG

    Dalam mode startup specific-offset, posisi dalam file log biner pada posisi startup.

    Untuk mendapatkan posisi log biner terbaru, jalankan pernyataan SHOW MASTER STATUS for table_name;.

    scan.startup.specific-offset.skip-events

    Tidak

    Tidak ada

    LONG

    Jumlah event yang dilewati setelah posisi startup yang ditentukan.

    scan.startup.specific-offset.skip-rows

    Tidak

    Tidak ada

    LONG

    Jumlah baris data yang dilewati setelah posisi startup yang ditentukan.

    scan.startup.timestamp-millis

    Tidak

    Tidak ada

    LONG

    Timestamp dalam milidetik dari posisi startup saat menggunakan mode startup timestamp.

    Saat menggunakan parameter ini, Anda harus mengatur scan.startup.mode ke timestamp. Timestamp dalam satuan milidetik (ms).

    server-time-zone

    Tidak

    Default sistem

    STRING

    Zona waktu sesi pada server database, seperti "Asia/Shanghai".

    Parameter ini mengontrol cara AnalyticDB for MySQLAnalyticDB for MySQL dikonversi ke tipe data STRING. Jika parameter ini tidak diatur, ZONELD.SYSTEMDEFAULT() digunakan untuk menentukan zona waktu server.

    debezium.min.row.count.to.stream.result

    Tidak

    1000

    INTEGER

    Jika jumlah baris dalam tabel lebih besar dari nilai ini, konektor melakukan streaming hasil.

    Jika Anda mengatur parameter ini ke 0, semua pemeriksaan ukuran tabel dilewati, dan semua hasil di-stream selama fase snapshot.

    connect.timeout

    Tidak

    30s

    DURATION

    Waktu maksimum menunggu koneksi database sebelum upaya tersebut timeout.

    Unit default adalah detik (s).

    connect.max-retries

    Tidak

    3

    INTEGER

    Jumlah maksimum percobaan ulang setelah kegagalan koneksi database.

  3. Buat tabel fisik di database tujuan untuk menyimpan data yang diproses. Topik ini menggunakan AnalyticDB for MySQL sebagai tujuan. Untuk konektor yang didukung oleh Flink, lihat Konektor yang didukung.

    CREATE TABLE target_table (
      `id` INT,
      `num` BIGINT,
      PRIMARY KEY (`id`)
    )
  4. Buat tabel sink yang terhubung ke tabel yang dibuat pada langkah sebelumnya. Tabel sink menulis data yang diproses ke tabel tertentu di AnalyticDB for MySQL.

    CREATE TEMPORARY TABLE adb_sink (
      `id` INT,
      `num` BIGINT,
      PRIMARY KEY (`id`) NOT ENFORCED
    ) WITH (
      'connector' = 'adb3.0',
      'url' = 'jdbc:mysql://amv-2zepb9n1l58ct01z50000****.ads.aliyuncs.com:3306/flinktest',
      'userName' = 'testUser',
      'password' = 'Test12****',
      'tableName' = 'target_table'
    );

    Untuk informasi lebih lanjut tentang parameter WITH dan pemetaan tipe untuk tabel sink, lihat Konektor AnalyticDB for MySQL V3.0.

  5. Gunakan pernyataan INSERT INTO untuk mengirim data dari tabel sumber ke tabel sink.

    INSERT INTO adb_sink
    SELECT * FROM adb_source;
  6. Klik Save.

  7. Klik Validate.

    Fitur validasi memeriksa semantik SQL, konektivitas jaringan, dan metadata tabel yang digunakan oleh pekerjaan. Anda juga dapat mengklik SQL Advice di area hasil untuk melihat prompt risiko SQL dan saran optimasi.

  8. (Opsional) Klik Debug.

    Anda dapat menggunakan fitur debugging pekerjaan untuk mensimulasikan eksekusi pekerjaan, memeriksa hasil output, memverifikasi logika bisnis pernyataan SELECT atau INSERT, meningkatkan efisiensi pengembangan, dan mengurangi risiko kualitas data.

  9. Klik Deploy.

    Setelah mengembangkan dan memvalidasi pekerjaan, terapkan ke lingkungan produksi. Lalu, buka halaman O&M dan mulai pekerjaan.

  10. (Opsional) Lihat informasi log biner.

    Catatan

    Pernyataan berikut mengembalikan 0 jika Anda telah mengaktifkan binary logging tetapi belum berlangganan menggunakan DTS. Informasi log biner hanya muncul setelah langganan berhasil dibuat.

    • Untuk mendapatkan nama file dan posisi entri log biner terbaru, jalankan pernyataan SQL berikut:

      SHOW MASTER STATUS FOR source_table;
    • Untuk melihat semua file log biner yang belum di-purge beserta ukurannya, jalankan pernyataan SQL berikut:

      SHOW BINARY LOGS FOR source_table;

Pemetaan tipe

Tabel berikut memetakan tipe data AnalyticDB for MySQL ke padanannya di Flink.

AnalyticDB for MySQL type

Flink type

BOOLEAN

BOOLEAN

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

DECIMAL(p,s) atau NUMERIC(p,s)

DECIMAL(p,s)

VARCHAR

STRING

BINARY

BYTES

DATE

DATE

TIME

TIME

DATETIME

TIMESTAMP

TIMESTAMP

TIMESTAMP

POINT

STRING

JSON

STRING