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.
CatatanUntuk 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:
-
Operasi DDL: Perubahan skema seperti ALTER TABLE tidak ditangkap. Jika skema tabel sumber berubah, Anda harus memperbarui definisi tabel secara manual dalam pekerjaan Flink dan menerapkan ulang pekerjaan tersebut.
-
INSERT OVERWRITE: Data yang ditulis melalui INSERT OVERWRITE tidak menghasilkan catatan log biner. Untuk mendapatkan data lengkap setelah operasi INSERT OVERWRITE, Anda perlu berlangganan ulang log biner.
-
Penghapusan partisi otomatis: Saat partisi yang kedaluwarsa dalam tabel partisi dihapus secara otomatis, data yang dihapus tidak menghasilkan catatan log biner.
-
Pembaruan penuh tampilan yang di-materialisasi: Operasi pembaruan penuh pada tampilan yang di-materialisasi tidak menghasilkan catatan log biner.
CatatanLog biner untuk tampilan yang di-materialisasi hanya berisi perubahan data dari pembaruan bertahap. Untuk mendapatkan status data lengkap setelah pembaruan penuh, Anda perlu berlangganan ulang log biner.
-
Langkah 1: Aktifkan binary logging
-
Aktifkan binary logging. Contoh ini menggunakan tabel bernama source_table.
CatatanAnalyticDB 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; -
(Opsional) Ubah periode retensi log biner.
Anda dapat mengubah parameter
binlog_ttluntuk 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_ttlmendukung format berikut:-
Milidetik: Nilai numerik. Contoh:
60merepresentasikan 60 milidetik. -
Detik: Angka + s. Contoh:
30smerepresentasikan 30 detik. -
Jam: Angka diikuti h. Contoh:
2hmerepresentasikan 2 jam. -
Hari: Angka diikuti
d. Contoh:1dmerepresentasikan 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
-
Unduh file JAR konektor: flink-sql-connector-adb-mysql-cdc-2.4-20260420.jar.
-
Masuk ke Konsol Realtime Compute for Apache Flink.
-
Pada tab Flink, temukan ruang kerja target dan klik Console di kolom Actions.
-
Di panel navigasi kiri, klik Connectors.
-
Pada halaman Connectors, klik Create Custom Connector.
-
Unggah konektor yang telah diunduh dan klik Next.
-
Klik Finish. Konektor kustom yang dibuat akan muncul dalam daftar konektor.
Langkah 3: Berlangganan log biner
-
Masuk ke Konsol Realtime Compute for Apache Flink dan buat pekerjaan SQL.
-
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.filedanscan.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.
PentingJika 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.modeke 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.
-
-
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`) ) -
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.
-
Gunakan pernyataan INSERT INTO untuk mengirim data dari tabel sumber ke tabel sink.
INSERT INTO adb_sink SELECT * FROM adb_source; -
Klik Save.
-
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.
-
(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.
-
Klik Deploy.
Setelah mengembangkan dan memvalidasi pekerjaan, terapkan ke lingkungan produksi. Lalu, buka halaman O&M dan mulai pekerjaan.
-
(Opsional) Lihat informasi log biner.
CatatanPernyataan 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 |