Realtime Compute for Apache Flink mendukung pernyataan CTAS (CREATE TABLE AS) untuk sinkronisasi satu tabel dan CDAS (CREATE DATABASE AS SELECT) untuk sinkronisasi multi-tabel atau seluruh database dari instans ApsaraDB RDS for MySQL ke kluster E-MapReduce (EMR) StarRocks secara real time.
Latar Belakang
Pernyataan CTAS (CREATE TABLE AS) secara otomatis membuat tabel StarRocks dengan skema yang identik dengan tabel sumber MySQL, serta menyinkronkan data dan mereplikasi perubahan skema dari sumber ke tujuan secara real time. Hal ini menyederhanakan pembuatan tabel tujuan dan menjaga konsistensi skema.
Saat mengeksekusi pernyataan CTAS, Flink melakukan langkah-langkah berikut:
-
Memeriksa apakah tabel tujuan sudah ada.
-
Jika tabel belum ada, Flink menggunakan katalog tujuan untuk membuat tabel tersebut dengan skema yang sama seperti tabel sumber.
-
Jika tabel sudah ada, Flink melewatkan pembuatan tabel dan melaporkan error jika skema tabel tujuan tidak sesuai dengan skema tabel sumber.
-
-
Mengirimkan dan memulai penerapan sinkronisasi data, termasuk sinkronisasi data historis dan perubahan skema dari tabel sumber ke tabel tujuan.
Pernyataan CTAS menyinkronkan data secara real time dan menyebarkan perubahan skema dari sumber ke tabel tujuan.
Perubahan skema mencakup pembuatan tabel awal dan perubahan struktur tabel berikutnya.
-
Perubahan skema yang didukung:
-
Menambahkan kolom nullable: Flink secara otomatis menambahkan kolom baru di akhir skema tabel tujuan dan menyinkronkan datanya.
-
Menghapus kolom nullable: Kolom tersebut tidak dihapus secara fisik dari tabel tujuan. Sebagai gantinya, Flink secara otomatis mengisi datanya dengan nilai
NULL. -
Mengganti nama kolom: Flink memperlakukan hal ini sebagai kombinasi penambahan kolom baru dan penghapusan kolom lama. Kolom yang diganti namanya ditambahkan di akhir tabel tujuan, dan data pada kolom asli diisi dengan nilai
NULL.Sebagai contoh, jika Anda mengganti nama
col_amenjadicol_b, kolomcol_bditambahkan di akhir tabel tujuan, dan data padacol_asecara otomatis diisi dengan nilaiNULL.
-
-
Perubahan skema yang tidak didukung:
-
Perubahan tipe data.
Contohnya, mengubah tipe dari
VARCHARkeBIGINTatau properti dariNOT NULLkeNULLABLE. -
Perubahan pada constraint, seperti kunci primer atau indeks.
-
Menambah atau menghapus kolom non-nullable.
-
Penyesuaian panjang field dalam pernyataan DDL.
-
-
Jika menghadapi perubahan skema yang tidak didukung, Anda harus menghapus secara manual tabel tujuan dan me-restart penerapan CTAS. Hal ini akan membuat ulang tabel tujuan dan menyinkronkan kembali seluruh data historis.
-
CTAS tidak mengenali jenis DDL tertentu. Sebaliknya, CTAS membandingkan perbedaan skema antara catatan data sebelum dan sesudah perubahan. Jika Anda menghapus kolom lalu menambahkannya kembali tanpa perubahan data di antara dua operasi DDL tersebut, CTAS tidak mendeteksi adanya perubahan skema. Demikian pula, jika Anda menambahkan kolom, CTAS hanya mendeteksi perubahan skema dan menyinkronkannya ke tabel tujuan setelah terjadi perubahan data pada tabel tersebut.
-
Untuk tipe data yang didukung saat membuat tabel dengan
CTAS, lihat Pemetaan tipe data antara Flink dan StarRocks. -
Saat menggunakan pernyataan CTAS untuk menggabungkan beberapa tabel MySQL, Flink secara otomatis menambahkan dua kolom,
_db_namedan_table_name, di awal skema tabel yang dihasilkan untuk melacak tabel sumber. Perilaku ini tidak dapat diubah. Oleh karena itu, saat menentukan urutan kolom untuk tabel baru, mulailah dari kolom ketiga agar skema hasil sesuai dengan ekspektasi Anda.
Prasyarat
-
Anda telah mengaktifkan Realtime Compute for Apache Flink (Fully Managed) dan membuat kluster Flink. Untuk informasi lebih lanjut, lihat Aktifkan Realtime Compute for Apache Flink (Fully Managed) dan Memulai penerapan Flink SQL.
-
Anda telah membuat kluster StarRocks. Untuk informasi lebih lanjut, lihat Buat kluster StarRocks.
-
Anda telah membuat instans ApsaraDB RDS for MySQL. Untuk informasi lebih lanjut, lihat Buat instans ApsaraDB RDS for MySQL.
Contoh dalam topik ini menggunakan MySQL 5.7, kluster E-MapReduce (EMR) StarRocks (EMR-3.39.1), dan Realtime Compute for Apache Flink (versi 1.15-vvr-6.0.3).
Batasan
-
Kluster Flink, kluster StarRocks, dan instans ApsaraDB RDS for MySQL harus berada dalam VPC yang sama.
-
Instans ApsaraDB RDS for MySQL harus versi 5.7 atau lebih baru.
-
Kluster StarRocks harus memiliki akses internet yang diaktifkan.
-
Versi Flink dalam kluster Flink Anda harus 1.15-vvr-6.0.3 atau lebih baru.
Langkah 1: Siapkan data uji
-
Buat database dan akun uji. Untuk informasi lebih lanjut, lihat Buat database dan akun untuk instans ApsaraDB RDS for MySQL.
Setelah membuat database dan akun, berikan izin baca dan tulis kepada akun uji tersebut.
CatatanDalam topik ini, database diberi nama
test_cdcdan akun diberi namatest. -
Gunakan akun uji untuk menghubungkan ke instans MySQL. Untuk informasi lebih lanjut, lihat Gunakan DMS untuk login ke instans ApsaraDB RDS for MySQL.
-
Jalankan perintah berikut di MySQL untuk membuat tabel data.
use test_cdc; CREATE TABLE IF NOT EXISTS `runoob_tbl`( `runoob_id` INT UNSIGNED AUTO_INCREMENT, `runoob_title` VARCHAR(100) NOT NULL, `runoob_author` VARCHAR(40) NOT NULL, `submission_date` DATE, `add_col` int DEFAULT NULL, PRIMARY KEY ( `runoob_id` ) )ENGINE=InnoDB DEFAULT CHARSET=utf8; INSERT INTO test_cdc.`runoob_tbl` (`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`) values (18,'first','tom','2022-06-22 17:13:44',3) -
Login ke kluster StarRocks menggunakan SSH. Untuk informasi lebih lanjut, lihat Login ke kluster.
-
Jalankan perintah berikut untuk menghubungkan ke kluster StarRocks.
mysql -h127.0.0.1 -P 9030 -uroot -
Jalankan perintah berikut untuk membuat pengguna dan memberikan izin.
CREATE DATABASE test_cdc; CREATE USER 'test' IDENTIFIED by '123456'; GRANT CREATE TABLE ON DATABASE test_cdc TO test;
Langkah 2: Buat katalog
Di halaman Draft Editor pada konsol Realtime Compute for Apache Flink, buat katalog untuk MySQL dan StarRocks. Untuk informasi lebih lanjut, lihat Memulai penerapan Flink SQL.
Parameter contoh hanya untuk referensi. Konfigurasikan sesuai kebutuhan Anda.
-
Katalog MySQL
-
Contoh
CREATE CATALOG mysql WITH ( 'type' = 'mysql', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'emr-test', 'password' = '123456', 'default-database' = 'test_cdc' ); -
Parameter
Parameter
Deskripsi
type
Tipe katalog. Tetapkan nilainya ke
mysql.hostname
Titik akhir internal instans ApsaraDB RDS for MySQL. Anda dapat menyalin titik akhir internal dari halaman Koneksi Database instans di konsol ApsaraDB RDS. Contohnya,
rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com.port
Nomor port layanan database MySQL. Nilai default-nya adalah
3306.username
Username untuk mengakses layanan database MySQL.
Gunakan username akun yang dibuat di Langkah 1: Siapkan data uji. Contoh ini menggunakan
test.password
Password untuk mengakses layanan database MySQL.
Gunakan password akun yang dibuat di Langkah 1: Siapkan data uji.
default-database
Nama database MySQL default.
Gunakan nama database dari Langkah 1: Siapkan data uji. Contoh ini menggunakan
test_cdc.
-
-
Katalog StarRocks
-
Contoh
CREATE CATALOG sr WITH ( 'type' = 'starrocks', 'endpoint' = '172.16.**.**:9030', 'username' = 'test', 'password' = '123456', 'dbname' = 'test_cdc' ); -
Parameter
Parameter
Deskripsi
type
Tipe katalog. Tetapkan nilainya ke
starrocks.endpoint
Alamat IP dan port StarRocks Frontend (FE).
username
Username untuk mengakses kluster StarRocks.
Gunakan username akun yang dibuat di Langkah 1: Siapkan data uji. Contoh ini menggunakan
test.password
Password untuk layanan database StarRocks.
Gunakan password akun yang dibuat di Langkah 1: Siapkan data uji.
dbname
Nama database StarRocks.
Gunakan nama database dari Langkah 1: Siapkan data uji. Contoh ini menggunakan
test_cdc.
-
Langkah 3: Buat dan publikasikan penerapan
-
Di halaman Draft Editor pada konsol Realtime Compute for Apache Flink, tulis pernyataan
CTAS.Berikut tiga contoh pernyataan CTAS.
-
Semantik at-least-once: Gunakan opsi sink.buffer-flush.interval-ms untuk mengonfigurasi interval penulisan data ke StarRocks. Opsi ini mengurangi latensi dan penggunaan memori.
/* At-least-once semantics */ use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl_sr with ( 'starrocks.create.table.properties'=' engine = olap primary key(runoob_id) distributed by hash(runoob_id ) buckets 8', 'database-name'='test_cdc', 'jdbc-url'='jdbc:mysql://172.16.**.**:9030', 'load-url'='172.16.**.**:18030', 'table-name'='runoob_tbl_sr', 'username'='test', 'password' = '123456', 'sink.buffer-flush.interval-ms' = '5000', 'sink.properties.row_delimiter' = '\x02', 'sink.properties.column_separator' = '\x01' ) as table mysql.test_cdc.runoob_tbl /*+ OPTIONS ( 'connector' = 'mysql-cdc', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'test', 'password' = '123456', 'database-name' = 'test_cdc', 'table-name' = 'runoob_tbl' )*/; -
Semantik exactly-once: Anda harus mengonfigurasi interval checkpoint. Hal ini mencegah kehilangan dan duplikasi data saat terjadi kegagalan, tetapi visibilitas data bergantung pada interval checkpoint. Untuk informasi lebih lanjut, lihat Checkpointing.
/* Exactly-once semantics. */ set 'execution.checkpointing.interval' = '1 min'; set 'execution.checkpointing.mode' = 'EXACTLY_ONCE'; set 'execution.checkpointing.timeout' = '10 min'; use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl with ( 'starrocks.create.table.properties'=' engine = olap primary key(runoob_id) distributed by hash(runoob_id ) buckets 8', 'database-name'='test_cdc', 'jdbc-url'='jdbc:mysql://172.16.**.**:9030', 'load-url'='172.16.**.**:18030', 'table-name'='runoob_tbl', 'username'='test', 'password' = '123456', 'sink.semantic' = 'exactly-once', 'sink.properties.row_delimiter' = '\x02', 'sink.properties.column_separator' = '\x01 ) as table mysql.test_cdc.runoob_tbl /*+ OPTIONS ( 'connector' = 'mysql-cdc', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'test', 'password' = '123456', 'database-name' = 'test_cdc', 'table-name' = 'runoob_tbl' )*/; -
Mode simple: Anda tidak perlu mendefinisikan field saat membuat tabel. Skema tabel disalin dari MySQL. Namun, Anda tidak dapat membuat partisi. Untuk menggunakan partisi, Anda harus menggunakan mode normal.
/* The two preceding examples use normal mode. This example demonstrates simple mode. */ use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl1 with ( 'starrocks.create.table.properties'='buckets 8', 'starrocks.create.table.mode'='simple', 'database-name'='test_cdc', 'jdbc-url'='jdbc:mysql://172.16.**.**:9030', 'load-url'='172.16.**.**:18030', 'table-name'='runoob_tbl_sr', 'username'='test', 'password' = '123456', 'sink.buffer-flush.interval-ms' = '5000', 'sink.properties.row_delimiter' = '\x02', 'sink.properties.column_separator' = '\x01' ) as table mysql.test_cdc.runoob_tbl /*+ OPTIONS ( 'connector' = 'mysql-cdc', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'emr-test', 'password' = '123456', 'database-name' = 'test_cdc', 'table-name' = 'runoob_tbl' )*/;
Tabel 1. Parameter WITH
Parameter
Wajib
Deskripsi
starrocks.create.table.properties
Ya
Definisi sufiks dalam pernyataan StarRocks
CREATE TABLE, tidak termasuk definisi kolom. Contohnya termasukengine,key, danbuckets.database-name
Ya
Nama database StarRocks.
Contoh ini menggunakan
test_cdc.jdbc-url
Ya
Digunakan untuk menjalankan operasi kueri di StarRocks.
Contohnya,
jdbc:mysql://172.16.**.**:9030. Bagian172.16.**.**adalah alamat IP internal kluster StarRocks.load-url
Ya
Alamat IP dan port HTTP StarRocks frontend (FE) dalam format
Alamat IP internal kluster StarRocks:Port. Topik ini menggunakan port 8030 sebagai contoh. Pilih port berdasarkan versi kluster Anda:18030: Untuk EMR V5.9.0 atau lebih baru dan EMR V3.43.0 atau lebih baru.
8030: Untuk EMR V5.8.0 atau lebih lama dan EMR V3.42.0 atau lebih lama.
CatatanUntuk informasi lebih lanjut tentang port, lihat UI dan port.
sink.semantic
Tidak
Tetapkan ke
exactly-onceuntuk menjamin semantik konsistensi data. Default-nya adalahat-least-once.starrocks.create.table.mode
Tidak
Nilai yang didukung:
-
normal(default): Anda harus memberikan konfigurasi lengkap, sepertiengine,key, danbuckets, dalam opsi starrocks.create.table.properties. -
simple: Engine diatur keolapdan tipe kunci diatur keprimary keysecara default. Kunci primer diwariskan dari tabel MySQL. Secara default, tabel didistribusikan berdasarkan hash semua kolom kunci primer, tanpa partisi. Anda harus menentukanbucketsdalam opsi starrocks.create.table.properties. Konfigurasi lain sepertipropertiesbersifat opsional.
Catatan-
Parameter
sink.use.new-apidihapus pada versi Flink 1.15-vvr-6.0.5 dan lebih baru. Jika Anda menggunakan versi sebelum 1.15-vvr-6.0.5, Anda harus menambahkan'sink.use.new-api'='false',ke parameter WITH. -
Untuk informasi tentang konfigurasi lainnya, lihat Muat data secara berkelanjutan dari Apache Flink.
Tabel 2. Parameter OPTIONS
Parameter
Deskripsi
connector
Tipe konektor. Tetapkan nilainya ke
mysql-cdc.hostname
Titik akhir internal instans ApsaraDB RDS for MySQL.
Anda dapat menyalin titik akhir internal dari halaman Koneksi Database instans di konsol ApsaraDB RDS. Contohnya,
rm-bp1nu0c46fn9k****.mysql.rds.aliyuncs.com.port
Nomor port layanan database MySQL. Nilai default-nya adalah
3306.username
Username untuk mengakses layanan database MySQL.
Gunakan username akun yang dibuat di Langkah 1: Siapkan data uji. Contoh ini menggunakan
test.password
Password untuk mengakses layanan database MySQL.
Gunakan password akun yang dibuat di Langkah 1: Siapkan data uji.
table-name
Nama tabel di StarRocks.
Gunakan nama tabel dari Langkah 1: Siapkan data uji. Contoh ini menggunakan
runoob_tbl.database-name
Nama database MySQL default.
Gunakan nama database dari Langkah 1: Siapkan data uji. Contoh ini menggunakan
test_cdc. -
-
Di pengaturan Lanjutan pada halaman Draft Editor, pilih versi Flink 1.15-vvr-6.0.3 atau lebih baru.
-
Klik online.
-
Di halaman Deployments, temukan penerapan target dan klik START di kolom Actions.
Langkah 4: Demonstrasi skenario
Kueri data
-
Login ke kluster StarRocks menggunakan SSH. Untuk informasi lebih lanjut, lihat Login ke kluster.
-
Jalankan perintah berikut untuk menghubungkan ke kluster StarRocks.
mysql -h127.0.0.1 -P 9030 -uroot -
Di CLI StarRocks, jalankan perintah berikut untuk melihat data tabel.
use test_cdc; select * from runoob_tbl1;Output menunjukkan bahwa data dari tabel MySQL telah disinkronkan ke StarRocks.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 18 | first | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Kueri data yang dimasukkan
-
Di konsol SQL untuk instans ApsaraDB RDS for MySQL, jalankan perintah berikut untuk memasukkan data.
INSERT INTO runoob_tbl(`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`) values(1,'second','tom2','2022-06-23',1) -
Di CLI StarRocks, jalankan perintah berikut untuk melihat data tabel.
select * from runoob_tbl1;Output menunjukkan bahwa data berhasil dimasukkan.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 1 | second | tom2 | 2022-06-23 | 1 | | 18 | first | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Sinkronisasi pembaruan data
-
Di konsol SQL untuk instans ApsaraDB RDS for MySQL, jalankan perintah berikut untuk memperbarui data tertentu.
update runoob_tbl set runoob_title= 'new' where runoob_id = 18 -
Di CLI StarRocks, jalankan perintah berikut untuk melihat data tabel.
select * from runoob_tbl1;Output menunjukkan bahwa pembaruan data telah disinkronkan.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 1 | second | tom2 | 2022-06-23 | 1 | | 18 | new | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Sinkronisasi penghapusan data
-
Di konsol SQL untuk instans ApsaraDB RDS for MySQL, jalankan perintah berikut untuk menghapus data tertentu.
DELETE FROM runoob_tbl WHERE runoob_id = 1 -
Di CLI StarRocks, jalankan perintah berikut untuk melihat data tabel.
select * from runoob_tbl1;Output menunjukkan bahwa penghapusan data telah disinkronkan.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 18 | new | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Menambahkan kolom nullable
-
Di konsol SQL untuk instans ApsaraDB RDS for MySQL, jalankan perintah berikut untuk menambahkan kolom nullable.
alter table `runoob_tbl` add COLUMN `add_col2` INT; -
Jalankan perintah berikut untuk memasukkan data.
INSERT INTO runoob_tbl(`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`,`add_col2`) values(1,'second','tom2','2022-06-23',1,2) -
Di CLI StarRocks, jalankan perintah berikut untuk melihat data tabel.
select * from runoob_tbl1;Output menunjukkan bahwa perubahan skema berhasil.
+-----------+--------------+---------------+-----------------+---------+----------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | add_col2 | +-----------+--------------+---------------+-----------------+---------+----------+ | 18 | new | tom | 2022-06-22 | 3 | NULL | +-----------+--------------+---------------+-----------------+---------+----------+ | 1 | second | tom2 | 2022-06-23 | 1 | 2 | | 18 | first | tom | 2022-06-22 | 3 | NULL | +-----------+--------------+---------------+-----------------+---------+----------+
Pengantar CDAS
Pernyataan CDAS merupakan syntactic sugar untuk CTAS. Pernyataan ini menyinkronkan seluruh database MySQL ke StarRocks sebagai satu penerapan Flink. Anda juga dapat menggunakan sintaks including table untuk menyinkronkan hanya subset tabel tertentu.
Seperti halnya CTAS, Anda harus membuat katalog MySQL dan StarRocks yang sesuai sebelum menjalankan pernyataan CDAS. Contoh berikut menunjukkan sintaksnya.
CREATE DATABASE IF NOT EXISTS sr_db with (
'starrocks.create.table.properties'=' buckets 8',
'starrocks.create.table.mode'='simple',
'jdbc-url'='jdbc:mysql://172.16.**.**:9030',
'load-url'='172.16.**.**:18030',
'username'='test',
'password' = '123456',
'sink.buffer-flush.interval-ms' = '5000' ,
'sink.properties.row_delimiter' = '\x02',
'sink.properties.column_separator' = '\x01'
)
as DATABASE mysql.test_cdc including table 'tabl1','tbl2','tbl3'
/*+ OPTIONS ( 'connector' = 'mysql-cdc',
'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com',
'port' = '3306',
'username' = 'test',
'password' = '123456',
'database-name' = 'test_cdc' )*/;