All Products
Search
Document Center

Realtime Compute for Apache Flink:Konektor MySQL YAML

Last Updated:Sep 18, 2026

Konektor MySQL dapat digunakan sebagai sumber data dalam pekerjaan Ingesti Data berbasis YAML.

Prasyarat

Sebelum menggunakan tabel sumber CDC MySQL, Anda harus menyelesaikan operasi prasyarat yang dijelaskan dalam Konfigurasi MySQL.

ApsaraDB RDS for MySQL

  • Lakukan probe jaringan untuk memastikan konektivitas jaringan ke Realtime Compute for Apache Flink.

  • Versi MySQL: 5.6, 5.7, 8.0.x, atau 8.4.

  • Pencatatan biner (binary logging) harus diaktifkan. Fitur ini diaktifkan secara default.

  • Format log biner harus ROW. Ini adalah format default.

  • Parameter `binlog_row_image` harus diatur ke FULL. Ini adalah pengaturan default.

  • Kompresi Transaksi Log Biner harus dinonaktifkan. Fitur ini diperkenalkan di MySQL 8.0.20 dan dinonaktifkan secara default.

  • Pengguna MySQL telah dibuat dengan izin SELECT, SHOW DATABASES, REPLICATION SLAVE, dan REPLICATION CLIENT.

  • Buat database dan tabel MySQL. Untuk informasi lebih lanjut, lihat Buat database dan akun untuk instans ApsaraDB RDS for MySQL. Gunakan akun istimewa untuk membuat database MySQL guna mencegah kegagalan operasi akibat izin yang tidak mencukupi.

  • Konfigurasikan daftar putih alamat IP. Untuk informasi lebih lanjut, lihat Konfigurasi daftar putih alamat IP untuk instans ApsaraDB RDS for MySQL.

PolarDB for MySQL

  • Lakukan probe jaringan untuk memastikan konektivitas jaringan ke Realtime Compute for Apache Flink.

  • Versi MySQL: 5.6, 5.7, 8.0.x, atau 8.4.

  • Pencatatan biner harus diaktifkan. Fitur ini dinonaktifkan secara default.

  • Format log biner harus ROW. Ini adalah format default.

  • Parameter `binlog_row_image` harus diatur ke FULL. Ini adalah pengaturan default.

  • Kompresi Transaksi Log Biner harus dinonaktifkan. Fitur ini diperkenalkan di MySQL 8.0.20 dan dinonaktifkan secara default.

  • Anda telah membuat pengguna MySQL dengan izin SELECT, SHOW DATABASES, REPLICATION SLAVE, dan REPLICATION CLIENT.

  • Buat database dan tabel MySQL. Untuk informasi lebih lanjut, lihat Buat database dan akun untuk kluster PolarDB for MySQL. Gunakan akun istimewa untuk membuat database MySQL guna mencegah kegagalan operasi akibat izin yang tidak mencukupi.

  • Konfigurasikan daftar putih alamat IP. Untuk informasi lebih lanjut, lihat Konfigurasi daftar putih alamat IP untuk kluster PolarDB for MySQL.

MySQL yang dikelola sendiri

  • Lakukan probe jaringan untuk memastikan konektivitas jaringan ke Realtime Compute for Apache Flink.

  • Versi MySQL: 5.6, 5.7, 8.0.x, atau 8.4.

  • Pencatatan biner harus diaktifkan. Fitur ini dinonaktifkan secara default.

  • Format log biner harus ROW. Format default adalah STATEMENT.

  • Parameter `binlog_row_image` harus diatur ke FULL. Ini adalah pengaturan default.

  • Kompresi Transaksi Log Biner harus dinonaktifkan. Fitur ini diperkenalkan di MySQL 8.0.20 dan dinonaktifkan secara default.

  • Buat pengguna MySQL dan berikan izin SELECT, SHOW DATABASES, REPLICATION SLAVE, dan REPLICATION CLIENT.

  • Buat database dan tabel MySQL. Untuk informasi lebih lanjut, lihat Buat database dan akun untuk instans MySQL yang dikelola sendiri. Gunakan akun istimewa untuk membuat database MySQL guna mencegah kegagalan operasi akibat izin yang tidak mencukupi.

  • Konfigurasikan daftar putih alamat IP. Untuk informasi lebih lanjut, lihat Konfigurasi daftar putih alamat IP untuk instans MySQL yang dikelola sendiri.

Batasan

Batasan umum

Konektor CDC MySQL tidak mendukung fitur Kompresi Transaksi Log Biner. Oleh karena itu, saat menggunakan konektor CDC MySQL untuk mengonsumsi data inkremental, pastikan Kompresi Transaksi Log Biner dinonaktifkan. Jika tidak, konektor mungkin gagal mengambil data inkremental.

ApsaraDB RDS for MySQL batasan

  • Untuk ApsaraDB RDS for MySQL, jangan membaca data dari database sekunder atau replika read-only. Hal ini karena periode retensi log biner default untuk database sekunder dan replika read-only sangat singkat. Jika log biner kedaluwarsa dan dihapus, pekerjaan mungkin gagal mengonsumsi data log biner dan melaporkan error.

  • ApsaraDB RDS for MySQL secara default mengaktifkan sinkronisasi primer/sekunder paralel tetapi tidak menjamin urutan transaksi yang konsisten antara instans primer dan sekunder. Hal ini dapat menyebabkan data terlewat selama alih bencana primer/sekunder dan pemulihan checkpoint. Untuk menghindari masalah ini, Anda dapat mengaktifkan opsi `slave_preserve_commit_order` secara manual untuk ApsaraDB RDS for MySQL.

PolarDB for MySQL batasan

Tabel sumber CDC MySQL tidak mendukung pembacaan data dari kluster Arsitektur Kluster Multi-master PolarDB for MySQL versi V1.0.19 dan sebelumnya. Untuk informasi lebih lanjut, lihat Apa itu Kluster Multi-master?. Log biner yang dihasilkan oleh kluster tersebut mungkin berisi ID tabel duplikat. Hal ini dapat menyebabkan error pemetaan skema pada tabel sumber CDC, yang mengakibatkan error saat mengurai data log biner.

MySQL open source batasan

Secara default, MySQL mempertahankan urutan transaksi selama replikasi log biner primer/sekunder. Jika replika MySQL memiliki replikasi paralel diaktifkan (slave_parallel_workers > 1) tetapi tidak memiliki slave_preserve_commit_order=ON diaktifkan, urutan commit transaksinya mungkin tidak konsisten dengan database primer. Saat Flink CDC melakukan pemulihan dari checkpoint, data mungkin terlewat karena urutan yang tidak sesuai. Anda dapat mengatur `slave_preserve_commit_order` = ON pada replika MySQL. Atau, Anda dapat mengatur `slave_parallel_workers` = 1, tetapi hal ini akan mengorbankan performa replikasi.

Catatan penggunaan

  • Tetapkan server ID untuk menghindari konflik konsumsi log biner.

  • Selama fase pembacaan data penuh, Anda tidak dapat menyimpan titik simpan (savepoint), menambahkan tabel baru ke atau menghapus tabel dari tabel sumber, lalu me-restart pekerjaan dari titik simpan tersebut. Hal ini akan menyebabkan pekerjaan gagal membaca data.

Ingesti Data

Anda dapat menggunakan konektor MySQL sebagai sumber data dalam pekerjaan Ingesti Data berbasis YAML.

Sintaks

source:
   type: mysql
   name: MySQL Source
   hostname: localhost
   port: 3306
   username: <username>
   password: <password>
   tables: adb.\.*, bdb.user_table_[0-9]+, [app|web].order_\.*
   server-id: 5401-5404

sink:
  type: xxx

Item konfigurasi

Parameter

Deskripsi

Wajib

Tipe data

Nilai default

Catatan

type

Jenis sumber data.

Ya

STRING

None

Nilainya harus mysql.

name

Nama sumber data.

Tidak

STRING

None

None.

hostname

Alamat IP atau hostname database MySQL.

Ya

STRING

None

Kami menyarankan agar Anda menentukan alamat VPC.

Catatan

Jika database MySQL dan Realtime Compute for Apache Flink tidak berada dalam VPC yang sama, Anda harus membuat koneksi jaringan cross-VPC atau menggunakan titik akhir publik untuk mengakses database. Untuk informasi lebih lanjut, lihat Kelola dan operasikan ruang kerja dan Bagaimana kluster Flink yang sepenuhnya dikelola mengakses Internet?.

username

Username untuk layanan database MySQL.

Ya

STRING

None

None.

password

Password untuk layanan database MySQL.

Ya

STRING

None

None.

tables

Tabel data MySQL yang akan disinkronkan.

Ya

STRING

None

  • Parameter ini mendukung ekspresi reguler untuk membaca data dari beberapa tabel.

  • Anda dapat menggunakan koma untuk memisahkan beberapa ekspresi reguler.

Catatan
  • Jangan gunakan karakter pencocokan awal string ^ dan akhir string $ dalam ekspresi reguler. Di VVR 11.2, titik digunakan untuk membagi ekspresi reguler guna mendapatkan bagian database. Karakter pencocokan awal dan akhir akan membuat ekspresi reguler database yang dihasilkan tidak dapat digunakan. Misalnya, Anda harus mengubah ^db.user_[0-9]+$ menjadi db.user_[0-9]+.

  • Titik digunakan untuk memisahkan nama database dan nama tabel. Untuk menggunakan titik sebagai pencocokan karakter apa pun, Anda harus meng-escape-nya dengan backslash. Contoh: db0.\.*, db1.user_table_[0-9]+, db[1-2].[app|web]order_\.*.

tables.exclude

Tabel yang akan dikecualikan dari sinkronisasi.

Tidak

STRING

None

  • Parameter ini mendukung ekspresi reguler untuk mengecualikan beberapa tabel.

  • Anda dapat menggunakan koma untuk memisahkan beberapa ekspresi reguler.

Catatan

Titik digunakan untuk memisahkan nama database dan nama tabel. Untuk menggunakan titik sebagai pencocokan karakter apa pun, Anda harus meng-escape-nya dengan backslash. Contoh: db0.\.*, db1.user_table_[0-9]+, db[1-2].[app|web]order_\.*.

port

Nomor port layanan database MySQL.

Tidak

INTEGER

3306

None.

schema-change.enabled

Menentukan apakah akan mengirimkan event perubahan skema.

Tidak

BOOLEAN

true

None.

server-id

ID numerik atau rentang untuk klien database yang digunakan untuk sinkronisasi.

Tidak

STRING

Nilai acak antara 5400 dan 6400 dihasilkan.

ID ini harus unik secara global dalam kluster MySQL. Tetapkan ID berbeda untuk setiap pekerjaan yang terhubung ke database yang sama. Parameter ini juga mendukung format rentang ID, seperti 5400-5408.

Catatan

Saat pembacaan inkremental diaktifkan, pembacaan konkuren didukung. Dalam kasus ini, kami menyarankan agar Anda menetapkan rentang ID sehingga setiap pembaca konkuren menggunakan ID yang berbeda.

jdbc.properties.*

Parameter koneksi kustom dalam URL JDBC.

Tidak

STRING

None

Anda dapat meneruskan parameter koneksi kustom. Misalnya, untuk tidak menggunakan protokol SSL, Anda dapat mengonfigurasi 'jdbc.properties.useSSL' = 'false'.

Untuk informasi lebih lanjut tentang parameter koneksi yang didukung, lihat MySQL Configuration Properties.

debezium.*

Parameter kustom untuk Debezium guna membaca log biner.

Tidak

STRING

None

Anda dapat meneruskan parameter Debezium kustom. Misalnya, gunakan 'debezium.event.deserialization.failure.handling.mode'='ignore' untuk menentukan logika penanganan error parsing.

Peringatan

Jangan mengubah parameter Debezium sembarangan. Hal ini dapat menyebabkan konektor salah membaca data. Misalnya, parameter debezium.binlog.buffer.size tidak boleh dikonfigurasi.

scan.incremental.snapshot.chunk.size

Ukuran setiap chunk dalam jumlah baris.

Tidak

INTEGER

8096

Tabel MySQL dibagi menjadi beberapa chunk untuk dibaca. Data sebuah chunk di-cache di memori sebelum sepenuhnya dibaca.

Semakin sedikit baris yang dimiliki setiap chunk, semakin besar jumlah total chunk dalam tabel. Meskipun hal ini mengurangi granularitas pemulihan kesalahan, namun dapat menyebabkan error OOM dan menurunkan throughput keseluruhan. Oleh karena itu, Anda perlu membuat pertimbangan dan menetapkan ukuran chunk yang wajar.

scan.snapshot.fetch.size

Jumlah maksimum catatan yang ditarik sekaligus saat membaca data lengkap sebuah tabel.

Tidak

INTEGER

1024

None.

scan.startup.mode

Mode startup untuk konsumsi data.

Tidak

STRING

initial

Nilai yang valid:

  • initial (default): Saat startup pertama kali atau startup tanpa status, konektor memindai data historis lengkap lalu membaca data log biner terbaru.

  • latest-offset: Saat startup pertama kali atau startup tanpa status, konektor tidak memindai data historis. Konektor mulai membaca dari akhir log biner, artinya hanya membaca perubahan terbaru setelah konektor dimulai.

  • earliest-offset: Konektor tidak memindai data historis. Konektor mulai membaca dari log biner paling awal yang tersedia.

  • specific-offset: Konektor tidak memindai data historis. Konektor mulai dari offset log biner tertentu. Anda dapat menentukan offset dengan mengonfigurasi scan.startup.specific-offset.file dan scan.startup.specific-offset.pos, atau hanya mengonfigurasi scan.startup.specific-offset.gtid-set untuk memulai dari set GTID tertentu.

  • timestamp: Konektor tidak memindai data historis. Konektor mulai membaca log biner dari timestamp tertentu. Timestamp ditentukan oleh scan.startup.timestamp-millis dalam milidetik.

Penting

Untuk mode startup earliest-offset, specific-offset, dan timestamp, jika skema tabel pada waktu startup berbeda dengan skema pada waktu offset awal yang ditentukan, pekerjaan akan melaporkan error karena ketidakcocokan skema. Dengan kata lain, saat menggunakan ketiga mode startup ini, Anda harus memastikan bahwa skema tabel yang bersangkutan tidak berubah antara posisi konsumsi log biner yang ditentukan dan waktu startup pekerjaan.

scan.startup.specific-offset.file

Nama file log biner untuk offset awal saat menggunakan mode startup specific-offset.

Tidak

STRING

None

Saat menggunakan parameter ini, Anda harus mengatur scan.startup.mode ke specific-offset. Contoh format nama file: mysql-bin.000003.

scan.startup.specific-offset.pos

Offset dalam file log biner yang ditentukan untuk offset awal saat menggunakan mode startup specific-offset.

Tidak

INTEGER

None

Saat menggunakan parameter ini, Anda harus mengatur scan.startup.mode ke specific-offset.

scan.startup.specific-offset.gtid-set

Set GTID untuk offset awal saat menggunakan mode startup specific-offset.

Tidak

STRING

None

Saat menggunakan parameter ini, Anda harus mengatur scan.startup.mode ke specific-offset. Contoh format set GTID: 24DA167-0C0C-11E8-8442-00059A3C7B00:1-19.

scan.startup.timestamp-millis

Timestamp dalam milidetik untuk offset awal saat menggunakan mode startup timestamp.

Tidak

LONG

None

Saat menggunakan parameter ini, Anda harus mengatur scan.startup.mode ke timestamp. Satuan timestamp adalah milidetik.

Penting

Saat Anda menentukan waktu, CDC MySQL mencoba membaca event awal setiap file log biner untuk menentukan timestamp-nya. Lalu, CDC MySQL menemukan file log biner yang sesuai dengan waktu yang ditentukan. Pastikan file log biner yang sesuai dengan timestamp yang ditentukan belum dihapus dari database dan dapat dibaca.

server-time-zone

Zona waktu sesi yang digunakan oleh database.

Tidak

STRING

Jika Anda tidak menentukan parameter ini, sistem menggunakan zona waktu lingkungan runtime pekerjaan Flink sebagai zona waktu server database. Ini adalah zona waktu yang Anda pilih.

Contoh: Asia/Shanghai. Parameter ini mengontrol bagaimana tipe TIMESTAMP di MySQL dikonversi ke tipe STRING. Untuk informasi lebih lanjut, lihat Debezium temporal values.

scan.startup.specific-offset.skip-events

Jumlah event log biner yang dilewati saat membaca dari offset tertentu.

Tidak

INTEGER

None

Saat menggunakan parameter ini, Anda harus mengatur scan.startup.mode ke specific-offset.

scan.startup.specific-offset.skip-rows

Jumlah perubahan baris yang dilewati saat membaca dari offset tertentu. Satu event log biner dapat berisi beberapa perubahan baris.

Tidak

INTEGER

None

Saat menggunakan parameter ini, Anda harus mengatur scan.startup.mode ke specific-offset.

connect.timeout

Waktu maksimum menunggu koneksi ke server database MySQL hingga timeout sebelum mencoba ulang.

Tidak

DURATION

30 s

None.

connect.max-retries

Jumlah maksimum percobaan ulang setelah koneksi ke layanan database MySQL gagal.

Tidak

INTEGER

3

None.

connection.pool.size

Ukuran kolam koneksi database.

Tidak

INTEGER

20

Kolam koneksi database digunakan untuk menggunakan kembali koneksi, yang dapat mengurangi jumlah koneksi database.

heartbeat.interval

Interval di mana sumber memajukan offset log biner menggunakan event heartbeat.

Tidak

DURATION

30s

Event heartbeat digunakan untuk memajukan offset log biner di sumber. Hal ini sangat berguna untuk tabel di MySQL yang jarang diperbarui. Untuk tabel seperti itu, offset log biner tidak dapat maju secara otomatis. Event heartbeat dapat mendorong offset log biner maju, yang mencegah masalah akibat offset log biner yang kedaluwarsa. Offset log biner yang kedaluwarsa dapat menyebabkan pekerjaan gagal dan tidak dapat dipulihkan, sehingga memerlukan restart tanpa status.

rds.region-id

ID wilayah instans Alibaba Cloud ApsaraDB RDS for MySQL.

Diperlukan saat menggunakan fitur untuk membaca log arsip dari OSS.

STRING

None

Untuk informasi lebih lanjut tentang ID wilayah, lihat Wilayah dan zona.

Penting

Karena string GTID untuk CDC MySQL dihasilkan secara acak dan tidak meningkat secara monoton seperti offset file log biner, menemukan GTID dalam file memerlukan pengunduhan dan parsing semua log arsip dari OSS. Proses ini sangat intensif sumber daya dan memakan waktu, sehingga fitur yang bergantung pada offset GTID tidak layak. Oleh karena itu, fitur log arsip OSS hanya mendukung pemulaian dari timestamp tertentu atau offset file log biner tertentu. Fitur ini tidak mendukung pemulaian dari GTID tertentu, maupun skenario dengan alih bencana primer/sekunder dalam log arsip, karena alih bencana primer/sekunder MySQL bergantung pada GTID. Evaluasi fitur ini secara hati-hati sebelum menggunakannya.

rds.access-key-id

ID AccessKey akun Alibaba Cloud ApsaraDB RDS for MySQL.

Diperlukan saat menggunakan fitur untuk membaca log arsip dari OSS.

STRING

None

Untuk informasi lebih lanjut, lihat Bagaimana cara melihat ID AccessKey dan Rahasia AccessKey?

Penting

Untuk mencegah Informasi AccessKey Anda bocor, gunakan fitur manajemen rahasia untuk menentukan ID AccessKey. Untuk informasi lebih lanjut, lihat Kelola variabel.

rds.access-key-secret

Rahasia AccessKey akun Alibaba Cloud ApsaraDB RDS for MySQL.

Diperlukan saat menggunakan fitur untuk membaca log arsip dari OSS.

STRING

None

Untuk informasi lebih lanjut, lihat Bagaimana cara melihat ID AccessKey dan Rahasia AccessKey?

Penting

Untuk mencegah Informasi AccessKey Anda bocor, gunakan fitur manajemen rahasia untuk menentukan Rahasia AccessKey. Untuk informasi lebih lanjut, lihat Kelola variabel.

rds.db-instance-id

ID instans Alibaba Cloud ApsaraDB RDS for MySQL.

Diperlukan saat menggunakan fitur untuk membaca log arsip dari OSS.

STRING

None

None.

rds.main-db-id

Nomor database utama instans Alibaba Cloud ApsaraDB RDS for MySQL.

Tidak

STRING

None

Untuk informasi lebih lanjut tentang cara mendapatkan nomor database utama, lihat Cadangan log ApsaraDB RDS for MySQL.

Catatan

Jika parameter ini tidak ditentukan, VVR 11.7 dan versi yang lebih baru secara otomatis mengambil nomor database utama berdasarkan informasi koneksi ApsaraDB RDS for MySQL.

rds.download.timeout

Periode timeout untuk mengunduh satu log arsip dari OSS.

Tidak

DURATION

60s

None.

rds.endpoint

Titik akhir layanan untuk mendapatkan informasi log biner OSS.

Tidak

STRING

None

Untuk informasi lebih lanjut tentang nilai yang valid, lihat Titik akhir.

rds.binlog-directory-prefix

Awalan direktori untuk menyimpan file log biner.

Tidak

STRING

rds-binlog-

None.

rds.use-intranet-link

Menentukan apakah akan menggunakan jaringan internal untuk mengunduh file log biner.

Tidak

BOOLEAN

true

None.

rds.binlog-directories-parent-path

Jalur mutlak direktori induk untuk menyimpan file log biner.

Tidak

STRING

None

None.

chunk-meta.group.size

Ukuran metadata chunk.

Tidak

INTEGER

1000

Jika metadata lebih besar dari nilai ini, metadata tersebut dibagi menjadi beberapa bagian untuk transmisi.

chunk-key.even-distribution.factor.lower-bound

Batas bawah faktor distribusi chunk untuk sharding merata.

Tidak

DOUBLE

0.05

Jika faktor distribusi kurang dari nilai ini, sharding tidak merata digunakan.

Faktor distribusi chunk = (MAX(chunk-key) - MIN(chunk-key) + 1) / Total jumlah baris data.

chunk-key.even-distribution.factor.upper-bound

Batas atas faktor distribusi chunk untuk sharding merata.

Tidak

DOUBLE

1000.0

Jika faktor distribusi lebih besar dari nilai ini, sharding tidak merata digunakan.

Faktor distribusi chunk = (MAX(chunk-key) - MIN(chunk-key) + 1) / Total jumlah baris data.

scan.incremental.close-idle-reader.enabled

Menentukan apakah akan menutup pembaca idle setelah snapshot selesai.

Tidak

BOOLEAN

false

Agar konfigurasi ini berlaku, Anda harus mengatur execution.checkpointing.checkpoints-after-tasks-finish.enabled ke true.

scan.only.deserialize.captured.tables.changelog.enabled

Pada fase inkremental, menentukan apakah hanya akan mendeserialisasi event perubahan dari tabel yang ditentukan.

Tidak

BOOLEAN

  • Nilai default adalah false di versi VVR 8.x.

  • Nilai default adalah true di VVR 11.1 dan versi yang lebih baru.

Nilai yang valid:

  • true: Hanya mendeserialisasi data perubahan dari tabel target untuk mempercepat pembacaan log biner.

  • false (default): Mendeserialisasi data perubahan dari semua tabel.

scan.parallel-deserialize-changelog.enabled

Pada fase inkremental, menentukan apakah akan menggunakan beberapa thread untuk mengurai event perubahan.

Tidak

BOOLEAN

false

Nilai yang valid:

  • true: Menggunakan beberapa thread pada fase deserialisasi event perubahan sambil mempertahankan urutan event log biner untuk mempercepat pembacaan.

  • false (default): Menggunakan satu thread pada fase deserialisasi event.

Catatan

Hanya didukung di VVR 8.0.11 dan versi yang lebih baru.

scan.parallel-deserialize-changelog.handler.size

Jumlah penanganan event saat menggunakan beberapa thread untuk mengurai event perubahan.

Tidak

INTEGER

2

Catatan

Hanya didukung di VVR 8.0.11 dan versi yang lebih baru.

metadata-column.include-list

Kolom metadata yang akan diteruskan ke downstream.

Tidak

STRING

None

Metadata yang tersedia meliputi op_ts, es_ts, query_log, file, dan pos. Anda dapat menggunakan koma untuk memisahkan beberapa kolom metadata.

Catatan

Konektor YAML CDC MySQL tidak memerlukan atau mendukung penambahan kolom metadata nama database, nama tabel, dan op_type. Anda dapat langsung menggunakan __data_event_type__ dalam ekspresi Transform untuk mendapatkan tipe data perubahan, atau menggunakan __schema_name__ dan __table_name__ untuk mendapatkan nama database dan nama tabel.

Penting
  • Kolom metadata file merepresentasikan file log biner tempat data berada. Kolom ini bernilai "" selama fase penuh dan nama file log biner selama fase inkremental. Kolom metadata pos merepresentasikan offset data dalam file log biner. Kolom ini bernilai "0" selama fase penuh dan offset data dalam file log biner selama fase inkremental. Kedua kolom metadata ini didukung mulai dari VVR 11.5.

  • Kolom metadata es_ts merepresentasikan waktu mulai transaksi yang sesuai untuk changelog di MySQL. Kolom ini hanya didukung untuk MySQL 8.0.x. Jangan tambahkan kolom metadata ini saat menggunakan versi MySQL sebelumnya.

  • Timestamp op_ts akurat hingga detik, sedangkan timestamp es_ts akurat hingga milidetik.

scan.newly-added-table.enabled

Saat me-restart dari checkpoint, menentukan apakah akan menyinkronkan tabel baru yang tidak cocok selama startup sebelumnya atau menghapus tabel dari state yang tidak lagi cocok.

Tidak

BOOLEAN

false

Ini berlaku saat me-restart dari checkpoint atau savepoint.

Penting

Selama fase pembacaan data penuh, Anda tidak dapat menyimpan savepoint, menambahkan tabel baru ke atau menghapus tabel dari tabel sumber, lalu me-restart pekerjaan dari savepoint tersebut. Hal ini akan menyebabkan pekerjaan gagal membaca data.

scan.binlog.newly-added-table.enabled

Pada fase inkremental, menentukan apakah akan mengirimkan data dari tabel baru yang cocok.

Tidak

BOOLEAN

false

Tidak dapat diaktifkan secara bersamaan dengan scan.newly-added-table.enabled.

scan.incremental.snapshot.chunk.key-column

Menentukan kolom untuk tabel tertentu yang akan digunakan sebagai kolom pemisah untuk sharding selama fase snapshot.

Tidak

STRING

None

  • Gunakan titik dua : untuk menghubungkan nama tabel dan nama kolom guna menentukan aturan. Nama tabel dapat berupa ekspresi reguler. Anda dapat menentukan beberapa aturan dengan memisahkannya menggunakan titik koma ;. Contoh: db1.user_table_[0-9]+:col1;db[1-2].[app|web]_order_\\.*:col2.

  • Wajib untuk tabel tanpa primary key. Kolom yang dipilih harus bertipe non-null (NOT NULL). Opsional untuk tabel dengan primary key. Hanya satu kolom yang dapat dipilih dari primary key.

scan.parse.online.schema.changes.enabled

Pada fase inkremental, menentukan apakah akan mencoba mengurai event DDL perubahan tanpa lock RDS.

Tidak

BOOLEAN

false

Nilai yang valid:

  • true: Mengurai event DDL perubahan tanpa lock RDS.

  • false (default): Tidak mengurai event DDL perubahan tanpa lock RDS.

Ini adalah fitur eksperimen. Sebelum melakukan perubahan tanpa lock online, ambil snapshot pekerjaan Flink untuk pemulihan.

Catatan

Hanya didukung di VVR 11.0 dan versi yang lebih baru.

scan.incremental.snapshot.backfill.skip

Menentukan apakah akan melewati backfill selama fase pembacaan snapshot.

Tidak

BOOLEAN

false

Nilai yang valid:

  • true: Melewati backfill selama fase pembacaan snapshot.

  • false (default): Tidak melewati backfill selama fase pembacaan snapshot.

Backfill hanya berlaku selama kueri snapshot satu chunk dan tidak mencakup seluruh fase pembacaan penuh. Saat backfill dilewati, kueri snapshot setiap chunk membaca data tabel terbaru pada saat itu; pembaruan yang terjadi pada chunk setelah dibaca tidak digabungkan selama fase pembacaan penuh dan dibaca dari Binlog setelah memasuki fase inkremental. Misalnya, pembaruan pada chunk5 yang terjadi saat chunk5 sedang di-snapshot tercermin langsung dalam snapshot chunk5; jika chunk5 diperbarui setelah pembaca maju ke chunk80, pembaruan tersebut diterapkan kemudian dari Binlog selama fase inkremental.

Penting

Saat diaktifkan, perubahan yang terjadi selama atau setelah pemindaian chunk tetap dikirimkan dari Binlog pada fase inkremental dan mungkin diduplikasi. Hanya semantik at-least-once yang dijamin. Aktifkan ini hanya jika sink downstream mendukung penulisan idempoten berdasarkan primary key.

Catatan

Hanya didukung di VVR 11.1 dan versi yang lebih baru.

treat-tinyint1-as-boolean.enabled

Menentukan apakah akan memperlakukan tipe TINYINT(1) sebagai tipe Boolean.

Tidak

BOOLEAN

true

Nilai yang valid:

  • true (default): Memperlakukan tipe TINYINT(1) sebagai tipe Boolean.

  • false: Tidak memperlakukan tipe TINYINT(1) sebagai tipe Boolean.

treat-timestamp-as-datetime-enabled

Menentukan apakah akan memperlakukan tipe TIMESTAMP sebagai tipe DATETIME.

Tidak

BOOLEAN

false

Nilai yang valid:

  • true: Memperlakukan tipe TIMESTAMP MySQL sebagai tipe DATETIME dan memetakannya ke tipe CDC TIMESTAMP.

  • false (default): Memetakan tipe TIMESTAMP MySQL ke tipe CDC TIMESTAMP_LTZ.

Tipe TIMESTAMP MySQL menyimpan waktu UTC dan dipengaruhi oleh zona waktu. Tipe DATETIME MySQL menyimpan waktu literal dan tidak dipengaruhi oleh zona waktu.

Saat diaktifkan, parameter ini mengonversi data tipe TIMESTAMP MySQL ke tipe DATETIME berdasarkan server-time-zone.

include-comments.enabled

Menentukan apakah akan menyinkronkan komentar tabel dan kolom.

Tidak

BOOELEAN

false

Nilai yang valid:

  • true: Menyinkronkan komentar tabel dan kolom.

  • false (default): Tidak menyinkronkan komentar tabel dan kolom.

Mengaktifkan opsi ini meningkatkan Penggunaan memori pekerjaan.

scan.incremental.snapshot.unbounded-chunk-first.enabled

Menentukan apakah akan mengirimkan chunk tak terbatas terlebih dahulu selama fase pembacaan snapshot.

Tidak

BOOELEAN

false

Nilai yang valid:

  • true: Mengirimkan chunk tak terbatas terlebih dahulu selama fase pembacaan snapshot.

  • false (default): Tidak mengirimkan chunk tak terbatas terlebih dahulu selama fase pembacaan snapshot.

Ini adalah fitur eksperimen. Mengaktifkannya dapat mengurangi risiko error OOM pada Pengelola Tugas saat menyinkronkan chunk terakhir selama fase snapshot. Tambahkan parameter ini sebelum startup pertama pekerjaan.

Catatan

Hanya didukung di VVR 11.1 dan versi yang lebih baru.

binlog.session.network.timeout

Timeout jaringan untuk koneksi log biner.

Tidak

DURATION

10m

Jika diatur ke 0s, timeout default server MySQL yang digunakan.

Catatan

Hanya didukung di VVR 11.5 dan versi yang lebih baru.

scan.rate-limit.records-per-second

Membatasi jumlah maksimum catatan yang dikirim oleh sumber per detik.

Tidak

LONG

None

Ini berlaku untuk skenario di mana pembacaan data perlu dibatasi. Batasan ini efektif baik pada fase penuh maupun inkremental.

Metrik numRecordsOutPerSecond dari sumber mencerminkan jumlah catatan yang di-output oleh seluruh aliran data per detik. Anda dapat menyesuaikan parameter ini berdasarkan metrik tersebut.

Pada fase pembacaan data penuh, biasanya Anda perlu mengurangi jumlah baris yang dibaca dalam setiap batch. Anda dapat mengurangi nilai parameter scan.incremental.snapshot.chunk.size.

Catatan

Hanya didukung di VVR 11.5 dan versi yang lebih baru.

include-binlog-meta.enable

Menentukan apakah akan menyertakan informasi log biner MySQL asli, seperti GTID dan offset log biner, dalam pesan.

Tidak

Boolean

false

Ini berlaku untuk skenario sinkronisasi log biner asli, seperti mengganti tautan sinkronisasi Canal yang sudah ada.

Catatan

Hanya didukung di VVR 11.6 dan versi yang lebih baru.

scan.binlog.tolerate.gtid-holes

Mengaktifkan parameter ini mengabaikan celah dalam urutan GTID, memungkinkan pekerjaan melewati event yang tidak kontinu dan terus berjalan.

Tidak

Boolean

false

Sebelum mengaktifkan parameter ini, Anda harus memastikan bahwa offset awal pekerjaan belum kedaluwarsa. Jika pekerjaan dimulai dari offset GTID yang telah dihapus atau kedaluwarsa, mesin akan diam-diam melewati log yang hilang, yang akan menyebabkan kehilangan data.

Catatan

Parameter ini hanya didukung di VVR 11.6 dan versi yang lebih baru.

scan.emit.create-table-events.in-batch.enabled

Menentukan apakah akan mengirimkan skema tabel secara batch selama fase inisialisasi pekerjaan.

Tidak

Boolean

false

Ini adalah fitur eksperimen. Aktifkan opsi ini saat satu pekerjaan menyinkronkan banyak tabel.

Catatan

Parameter ini hanya didukung di VVR 11.4 dan versi yang lebih baru.

Gunakan katalog yang sudah ada

Mulai dari VVR 11.5, Anda dapat langsung mereferensikan katalog MySQL bawaan yang dibuat di halaman Data Management dalam pekerjaan Ingesti Data Flink CDC. Hal ini mengurangi upaya manual dalam menulis properti koneksi.

source:
  type: mysql
  using.built-in-catalog: mysql_rds_catalog

Saat ini, pekerjaan Ingesti Data mendukung penggunaan ulang otomatis parameter katalog MySQL berikut:

  • hostname

  • port

  • username

  • password

  • catalog.table.metadata-columns

  • catalog.table.treat-tinyint1-as-boolean

Jika Anda ingin mengganti salah satu parameter yang digunakan ulang secara otomatis ini, Anda dapat menulis eksplisit parameter YAML yang sesuai. Parameter yang ditulis eksplisit memiliki prioritas lebih tinggi.

Pemetaan tipe

Tabel berikut menunjukkan pemetaan tipe data untuk Ingesti Data.

Tipe bidang CDC MySQL

Tipe bidang CDC

TINYINT(n)

TINYINT

SMALLINT

SMALLINT

TINYINT UNSIGNED

TINYINT UNSIGNED ZEROFILL

YEAR

INT

INT

MEDIUMINT

MEDIUMINT UNSIGNED

MEDIUMINT UNSIGNED ZEROFILL

SMALLINT UNSIGNED

SMALLINT UNSIGNED ZEROFILL

BIGINT

BIGINT

INT UNSIGNED

INT UNSIGNED ZEROFILL

BIGINT UNSIGNED

DECIMAL(20, 0)

BIGINT UNSIGNED ZEROFILL

SERIAL

FLOAT [UNSIGNED] [ZEROFILL]

FLOAT

DOUBLE [UNSIGNED] [ZEROFILL]

DOUBLE

DOUBLE PRECISION [UNSIGNED] [ZEROFILL]

REAL [UNSIGNED] [ZEROFILL]

NUMERIC(p, s) [UNSIGNED] [ZEROFILL] dan p <= 38

DECIMAL(p, s)

DECIMAL(p, s) [UNSIGNED] [ZEROFILL] dan p <= 38

FIXED(p, s) [UNSIGNED] [ZEROFILL] dan p <= 38

BOOLEAN

BOOLEAN

BIT(1)

TINYINT(1)

DATE

DATE

TIME [(p)]

TIME [(p)]

DATETIME [(p)]

TIMESTAMP [(p)]

TIMESTAMP [(p)]

Pemetaan tergantung pada nilai parameter treat-timestamp-as-datetime-enabled:

true:TIMESTAMP[(p)]

false:TIMESTAMP_LTZ[(p)]

CHAR(n)

CHAR(n)

VARCHAR(n)

VARCHAR(n)

BIT(n)

BINARY(⌈(n + 7) / 8⌉)

BINARY(n)

BINARY(n)

VARBINARY(N)

VARBINARY(N)

NUMERIC(p, s) [UNSIGNED] [ZEROFILL] dan 38 < p <= 65

STRING

Catatan

Di MySQL, tipe data desimal memiliki presisi hingga 65, tetapi di Flink, presisi dibatasi hingga 38. Oleh karena itu, jika Anda mendefinisikan kolom desimal dengan presisi lebih dari 38, Anda harus memetakannya ke string untuk menghindari kehilangan presisi.

DECIMAL(p, s) [UNSIGNED] [ZEROFILL] dan 38 < p <= 65

FIXED(p, s) [UNSIGNED] [ZEROFILL] dan 38 < p <= 65

TINYTEXT

STRING

TEXT

MEDIUMTEXT

LONGTEXT

ENUM

JSON

STRING

Catatan

Tipe data JSON dikonversi ke string berformat JSON di Flink.

GEOMETRY

STRING

Catatan

Tipe data spasial di MySQL dikonversi ke string dengan format JSON tetap. Untuk informasi lebih lanjut, lihat Pemetaan Tipe Data Spasial MySQL dari MySQL.

POINT

LINESTRING

POLYGON

MULTIPOINT

MULTILINESTRING

MULTIPOLYGON

GEOMETRYCOLLECTION

TINYBLOB

BYTES

Catatan

Untuk tipe data BLOB di MySQL, hanya blob dengan panjang tidak lebih dari 2.147.483.647 (2**31-1) yang didukung.

BLOB

MEDIUMBLOB

LONGBLOB

Tetapkan server ID untuk menghindari konflik konsumsi log biner

Saat pekerjaan Ingesti Data membaca log biner, sumber mendaftar ke MySQL sebagai klien replikasi menggunakan server-id. Jika beberapa pekerjaan atau klien replikasi lain menggunakan server ID yang sama, terjadi konflik konsumsi log biner dan pekerjaan gagal. Perhatikan hal berikut saat mengonfigurasi server ID:

  • Secara default, server-id adalah nilai tunggal acak antara 5400 dan 6400. Jika beberapa pekerjaan menggunakan nilai default, konflik mungkin terjadi. Kami menyarankan agar Anda secara eksplisit mengonfigurasi ID yang tidak tumpang tindih untuk pekerjaan tersebut.

  • Jika paralelisme sumber lebih dari 1, Anda harus mengonfigurasi rentang server ID. Jumlah ID yang tersedia dalam rentang tersebut harus tidak kurang dari paralelisme, dan setiap pembaca paralel menggunakan ID yang berbeda.

Contoh: Jika paralelisme sumber adalah 4, konfigurasikan rentang yang berisi empat ID.

source:
  type: mysql
  name: MySQL Source
  hostname: <hostname>
  port: 3306
  username: <username>
  password: <password>
  tables: app_db.\.*
  server-id: 5400-5403

sink:
  type: hologres

Jika beberapa pekerjaan membaca dari instans MySQL yang sama, tetapkan rentang yang tidak tumpang tindih untuk setiap pekerjaan. Misalnya, pekerjaan A menggunakan 5400-5403, dan pekerjaan B menggunakan 5404-5407.

Percepat pembacaan log biner

Saat Anda menggunakan konektor MySQL sebagai sumber data Ingesti, konektor mengurai file log biner untuk menghasilkan berbagai pesan perubahan selama fase inkremental. File log biner mencatat semua perubahan tabel dalam format biner. Anda dapat mempercepat penguraian file log biner dengan cara berikut.

  • Aktifkan parsing paralel dan filter parsing (Fitur ini memerlukan Realtime Compute for Apache Flink dengan Ververica Runtime (VVR) 8.0.7 atau versi yang lebih baru. Fitur ini tidak tersedia di edisi komunitas konektor CDC MySQL.)

    • Aktifkan opsi scan.only.deserialize.captured.tables.changelog.enabled untuk mengurai hanya event perubahan dari tabel yang ditentukan.

    • Aktifkan opsi scan.parallel-deserialize-changelog.enabled untuk menggunakan beberapa thread mengurai file log biner dan mengirimkan event ke antrian konsumen secara berurutan. Saat Anda mengaktifkan opsi ini, biasanya Anda juga perlu meningkatkan CPU TaskManager.

  • Optimalkan parameter Debezium

    debezium.max.queue.size: 162580
    debezium.max.batch.size: 40960
    debezium.poll.interval.ms: 50
    • debezium.max.queue.size: Jumlah maksimum catatan yang dapat ditampung oleh antrian pemblokiran. Saat Debezium membaca aliran event dari database, event tersebut ditempatkan dalam antrian pemblokiran sebelum ditulis ke downstream. Nilai default adalah 8192.

    • debezium.max.batch.size: Jumlah maksimum event yang diproses konektor dalam setiap iterasi. Nilai default adalah 2048.

    • debezium.poll.interval.ms: Jumlah milidetik yang harus ditunggu konektor sebelum meminta event perubahan baru. Nilai default adalah 1000 milidetik, atau 1 detik.

Contoh penggunaan:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: ${mysql.source.table}
  server-id: 7601-7604
  # Konfigurasi Debezium
  debezium.max.queue.size: 162580
  debezium.max.batch.size: 40960
  debezium.poll.interval.ms: 50
  # Aktifkan filter parsing
  scan.only.deserialize.captured.tables.changelog.enabled: true

Kapasitas konsumsi log biner Edisi Perusahaan CDC MySQL adalah 85 MB/detik, sekitar dua kali lipat dari versi komunitas open source. Saat kecepatan pembuatan file log biner melebihi 85 MB/detik (yaitu, satu file 512 MB setiap 6 detik), latensi pekerjaan Flink terus meningkat. Latensi pemrosesan secara bertahap berkurang setelah kecepatan pembuatan file log biner melambat. Jika file log biner berisi transaksi besar, latensi pemrosesan mungkin meningkat sementara. Latensi berkurang setelah log untuk transaksi tersebut dibaca.

Diagnosis latensi data untuk mengoptimalkan throughput pekerjaan

Jika Anda mengalami latensi data selama fase inkremental, analisis masalah dengan mengikuti langkah-langkah berikut:

  1. Periksa metrik currentFetchEventTimeLag dan currentEmitEventTimeLag di halaman Ikhtisar. Metrik currentFetchEventTimeLag merepresentasikan latensi dalam membaca data dari log biner. Metrik currentEmitEventTimeLag merepresentasikan latensi dalam membaca data untuk tabel yang relevan dengan pekerjaan dari log biner.

    Skenario

    Deskripsi

    currentFetchEventTimeLag rendah, sedangkan currentEmitEventTimeLag tinggi dan jarang diperbarui.

    currentFetchEventTimeLag yang rendah menunjukkan bahwa menarik log biner dari database efisien. Namun, log biner berisi sedikit data untuk tabel yang perlu dibaca pekerjaan. Oleh karena itu, currentEmitEventTimeLag jarang diperbarui. Ini adalah perilaku yang diharapkan.

    Kedua metrik currentFetchEventTimeLag dan currentEmitEventTimeLag tinggi.

    Hal ini menunjukkan bahwa tabel sumber memiliki performa baca yang buruk. Anda dapat melanjutkan ke langkah-langkah berikutnya di bagian ini untuk optimasi.

  2. Backpressure dapat mengurangi laju pengiriman data oleh sumber ke operator downstream. Anda mungkin mengamati bahwa sourceIdleTime meningkat secara periodik, dan kedua metrik currentFetchEventTimeLag serta currentEmitEventTimeLag terus bertambah. Untuk mengatasi hal ini, tingkatkan paralelisme node tempat backpressure berasal.

  3. Periksa metrik TM CPU Usage di halaman CPU dan metrik TM GC Time di halaman JVM untuk menentukan apakah sumber daya CPU atau memori tidak mencukupi. Anda dapat meningkatkan sumber daya pekerjaan untuk mengoptimalkan performa baca.

Baca log biner arsip dari OSS

Saat Anda menggunakan instans ApsaraDB RDS for MySQL sebagai sumber data, Anda dapat membaca cadangan log yang disimpan di OSS. Jika file yang sesuai dengan timestamp atau posisi log biner yang ditentukan disimpan di OSS, Flink secara otomatis menarik file log dari OSS ke kluster. Jika file disimpan secara lokal di database, Flink secara otomatis beralih ke pembacaan melalui koneksi database. Fitur ini hanya tersedia di Realtime Compute for Apache Flink dan tidak didukung di edisi komunitas konektor CDC MySQL.

Untuk mengaktifkan pembacaan dari cadangan log OSS, Anda harus mengonfigurasi parameter koneksi ApsaraDB RDS for MySQL. Contoh:

source:
  type: mysql
  hostname: <yourHostname>
  port: 3306
  username: <yourUsername>
  password: <yourPassword>
  tables: <yourTables>
  # Parameter koneksi RDS untuk membaca log biner arsip dari OSS
  rds.region-id: cn-beijing
  rds.access-key-id: your_access_key_id
  rds.access-key-secret: your_access_key_secret
  rds.db-instance-id: rm-xxxxxxxx # ID instans database.
  rds.main-db-id: 12345678 # ID database utama.
  rds.endpoint: rds.aliyuncs.com

Ingesti data untuk sinkronisasi database dan skema

Untuk pekerjaan yang hanya melibatkan logika sinkronisasi data, kami menyarankan menjalankannya sebagai pekerjaan Ingesti Data. Pekerjaan Ingesti Data dioptimalkan secara mendalam untuk skenario integrasi data. Untuk petunjuk penggunaan, lihat Memulai Ingesti Data Flink CDC dan Mengembangkan Pekerjaan Ingesti Data Flink CDC (Pratinjau Publik).

Kode berikut menunjukkan contoh pekerjaan Ingesti Data yang menyinkronkan seluruh database app_db dari MySQL ke Hologres, termasuk perubahan skema selanjutnya dari database hulu:

source:
  type: mysql
  hostname: <hostname>
  port: 3306
  username: ${secret_values.mysqlusername}
  password: ${secret_values.mysqlpassword}
  tables: app_db.\.*
  server-id: 5400-5404

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

Penemuan tabel baru dalam ingesti data

Konektor Ingesti Data MySQL menyediakan opsi konfigurasi untuk mendukung penemuan tabel baru dalam dua skenario berbeda.

Parameter

Deskripsi

Catatan

scan.newly-added-table.enabled

Saat pekerjaan me-restart dari checkpoint, opsi ini menyinkronkan tabel yang tidak ditemukan selama startup sebelumnya. Opsi ini membaca data snapshot dan inkremental dari tabel baru tersebut.

Opsi ini hanya didukung saat scan.startup.mode diatur ke initial. Opsi ini tidak berpengaruh pada mode startup lainnya.

scan.binlog.newly-added-table.enabled

Selama fase inkremental, opsi ini secara otomatis menyinkronkan data dari tabel baru yang ditemukan.

  • Kami menyarankan mengaktifkan opsi ini saat pertama kali memulai pekerjaan. Pekerjaan secara otomatis mengurai pernyataan DDL CREATE TABLE dan menyinkronkan data ke downstream. Jika Anda mengaktifkan opsi ini dan me-restart pekerjaan setelah tabel database sudah dibuat, hal ini dapat menyebabkan data tidak lengkap.

  • Dalam mode startup initial, operasi DDL tidak disinkronkan ke downstream hingga fase snapshot selesai. Tabel yang dibuat selama fase snapshot tidak dapat disinkronkan secara otomatis meskipun scan.binlog.newly-added-table.enabled diaktifkan.

Penting
  • Selama fase pembacaan penuh, menyimpan savepoint lalu menambahkan atau menghapus tabel sumber sebelum me-restart dari savepoint tersebut tidak didukung. Melakukan hal ini mencegah pekerjaan membaca data dengan benar.

  • Jangan mengaktifkan scan.newly-added-table.enabled dan scan.binlog.newly-added-table.enabled secara bersamaan. Mengaktifkan keduanya dapat menyebabkan duplikasi data.