Konektor AnalyticDB for MySQL V3.0 memungkinkan Anda membaca dari dan menulis ke kluster AnalyticDB for MySQL V3.0 menggunakan Flink SQL. AnalyticDB for MySQL adalah layanan data warehousing cloud-native yang mendukung penulisan real-time ber-throughput tinggi, analitik latensi rendah, serta operasi ekstrak, transformasi, dan muat (ETL) yang kompleks.
Konektor ini mendukung jenis tabel dan kemampuan berikut:
| Item | Deskripsi |
|---|---|
| Jenis tabel | Tabel sumber, tabel dimensi, dan Tabel sink. Catatan
Tabel sumber memerlukan Ververica Runtime (VVR) 8.0.4 atau versi yang lebih baru. Untuk parameter tabel sumber, lihat Gunakan Flink untuk berlangganan Log biner. |
| Mode eksekusi | Streaming mode dan batch mode |
| Format data | N/A |
| Metrik | N/A |
| Jenis API | SQL API |
| Pembaruan dan penghapusan data pada tabel sink | Didukung |
Prasyarat
Sebelum memulai, pastikan Anda telah:
-
Membuat kluster AnalyticDB for MySQL dan sebuah tabel. Lihat Buat kluster dan CREATE TABLE.
-
Mengonfigurasi Daftar putih alamat IP untuk kluster tersebut. Lihat Konfigurasikan Daftar putih alamat IP.
Sintaksis
CREATE TEMPORARY TABLE adb_table (
`id` INT,
`num` BIGINT,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'adb3.0',
'url' = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
'userName' = '<yourUsername>',
'password' = '<yourPassword>',
'tableName' = '<yourTablename>'
);
Kunci primary yang didefinisikan dalam pernyataan DDL Flink harus sesuai dengan kunci primary tabel fisik di AnalyticDB for MySQL—bidang yang sama, nama yang sama, dan harus ada di kedua sisi. Ketidaksesuaian kunci primary dapat menyebabkan korupsi data.
Perilaku utama
Kunci primary dan mode penulisan
Jika kunci primary didefinisikan dalam DDL, konektor beroperasi dalam mode upsert, di mana kunci primary duplikat akan memicu pembaruan alih-alih penyisipan baru. Perilaku SQL yang tepat bergantung pada parameter replaceMode:
-
replace— menggunakanREPLACE INTO. Kunci primary duplikat akan menimpa seluruh baris yang sudah ada. -
upsert— menggunakanINSERT INTO ... ON DUPLICATE KEY UPDATE. Hanya bidang yang ditentukan yang diperbarui; bidang lain tetap mempertahankan nilai saat ini. -
insert— menggunakanINSERT IGNORE INTO. Kunci primary duplikat diabaikan secara diam-diam; baris yang sudah ada tetap dipertahankan.
Jika tidak ada kunci primary yang didefinisikan, konektor selalu menggunakan INSERT IGNORE INTO.
replaceModememerlukan AnalyticDB for MySQL V3.1.3.5 atau versi yang lebih baru serta VVR 11.2 atau versi yang lebih baru untuk nilai string (replace,upsert,insert). Versi VVR sebelumnya menggunakantrue(setara denganreplace) danfalse(setara denganupsert). VVR 11.2 dan versi setelahnya tetap kompatibel dengantruedanfalse.
Pembuferan penulisan
Konektor menyimpan catatan dalam memori dan mengirimkannya secara batch. Pengiriman dilakukan ketika salah satu kondisi berikut terpenuhi:
-
Jumlah catatan yang dibuffer mencapai
batchSize(default: 1.000) ataubufferSize(default: 1.000). -
Waktu sejak pengiriman terakhir mencapai
flushIntervalMs(default: 3.000 ms).
Baik batchSize maupun bufferSize hanya berlaku jika kunci primary didefinisikan.
Peng-cache-an tabel dimensi
Konektor mendukung tiga kebijakan cache untuk tabel dimensi. Penggunaan cache mengurangi jumlah lookup ke tabel fisik tetapi memerlukan memori tambahan.
| Kebijakan cache | Perilaku | Gunakan saat |
|---|---|---|
ALL (default) |
Seluruh data tabel dimensi dimuat ke memori sebelum pekerjaan dimulai. Lookup selanjutnya hanya mengakses cache. Data dimuat ulang setelah cacheTTLMs kedaluwarsa. |
Tabel dimensi berukuran kecil dan kunci yang tidak ditemukan sering terjadi. |
LRU |
Baris yang sering diakses di-cache. Saat terjadi cache miss, konektor melakukan query ke tabel fisik dan memperbarui cache. Entri kedaluwarsa setelah cacheTTLMs. |
Tabel dimensi berukuran besar tetapi aksesnya condong ke subset kunci tertentu. |
None |
Tanpa caching. Setiap lookup langsung mengakses tabel fisik. | Volume lookup rendah atau persyaratan kesegaran data sangat ketat. |
Dengan caching ALL, konektor memuat seluruh tabel dimensi ke memori secara asinkron. Anda harus meningkatkan memori node yang digunakan untuk join tabel—minimal dua kali ukuran tabel remote—untuk menghindari error out of memory (OOM). Pantau Penggunaan memori node selama pekerjaan berjalan.
Parameter dalam klausa WITH
Parameter umum
| Parameter | Tipe | Wajib | Default | Deskripsi |
|---|---|---|---|---|
connector |
String | Ya | — | Diatur ke adb3.0 |
url |
String | Ya | — | URL Java Database Connectivity (JDBC) database, dalam format jdbc:mysql://<endpoint>:<port>/<databaseName>. Dapatkan endpoint dan Port dari bagian Network Information pada halaman kluster Anda di Konsol AnalyticDB for MySQL. |
userName |
String | Ya | — | Username untuk akses database |
password |
String | Ya | — | Password untuk akses database |
tableName |
String | Ya | — | Nama tabel target di database |
maxRetryTimes |
Integer | Tidak | 10 | Jumlah maksimum percobaan ulang saat terjadi kegagalan baca atau tulis |
Parameter tabel sink
| Parameter | Tipe | Wajib | Default | Deskripsi |
|---|---|---|---|---|
batchSize |
Integer | Tidak | 1000 | Jumlah catatan yang ditulis per batch. Hanya berlaku jika kunci primary didefinisikan. |
bufferSize |
Integer | Tidak | 1000 | Jumlah maksimum catatan yang dibuffer di memori sebelum pemicu flush. Hanya berlaku jika kunci primary didefinisikan. |
flushIntervalMs |
Integer | Tidak | 3000 | Waktu maksimum (dalam milidetik) antar flush. Ketika interval ini berakhir, semua catatan yang dibuffer akan ditulis terlepas dari ukuran buffer. |
ignoreDelete |
Boolean | Tidak | false | Apakah operasi penghapusan diabaikan. Atur ke true untuk melewati penghapusan; false untuk menerapkannya. |
replaceMode |
Boolean | Tidak | true |
Mode penulisan saat kunci primary didefinisikan. VVR 11.2+: replace, upsert, atau insert. Versi sebelumnya: true (sama dengan replace) atau false (sama dengan upsert). Memerlukan AnalyticDB for MySQL V3.1.3.5 atau versi yang lebih baru. Hanya berlaku jika kunci primary didefinisikan dalam DDL; jika tidak, INSERT IGNORE INTO selalu digunakan. |
excludeUpdateColumns |
String | Tidak | (kosong) | Daftar kolom yang dipisahkan koma untuk dikecualikan dari pembaruan saat replaceMode bernilai upsert (atau false). Saat terjadi kunci primary duplikat: hanya kolom yang tersisa yang diperbarui; kolom yang dikecualikan mempertahankan nilai yang sudah ada. Saat kunci primary unik: semua kolom disisipkan. Contoh: excludeUpdateColumns='column1,column2'. Nilainya harus berada dalam satu baris tanpa jeda baris. |
connectionMaxActive |
Integer | Tidak | 40 | Jumlah maksimum koneksi database konkuren dalam kolam thread |
Parameter tabel dimensi
| Parameter | Tipe | Wajib | Default | Deskripsi |
|---|---|---|---|---|
cache |
String | Tidak | ALL | Kebijakan cache: None, LRU, atau ALL. Lihat Peng-cache-an tabel dimensi. |
cacheSize |
Integer | Tidak | 100000 | Jumlah maksimum baris yang di-cache. Wajib saat cache bernilai LRU. |
cacheTTLMs |
Integer | Tidak | Long.MAX_VALUE | Masa berlaku entri cache dalam milidetik. Untuk LRU: entri kedaluwarsa setelah periode ini (default: tidak pernah kedaluwarsa). Untuk ALL: seluruh cache dimuat ulang setelah periode ini (default: tidak pernah dimuat ulang). Tidak digunakan saat cache bernilai None. |
maxJoinRows |
Integer | Tidak | 1024 | Jumlah maksimum baris tabel dimensi yang cocok per catatan input. Jika setiap catatan input cocok dengan paling banyak n baris tabel dimensi, atur nilai ini ke n untuk mengoptimalkan performa join di Realtime Compute for Apache Flink. |
Pemetaan tipe data
| AnalyticDB for MySQL V3.0 | Realtime Compute for Apache Flink |
|---|---|
| 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 |
Contoh
Tabel sink
Contoh berikut membaca dari sumber datagen dan menulis data ke tabel sink AnalyticDB for MySQL.
CREATE TEMPORARY TABLE datagen_source (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE adb_sink (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'adb3.0',
'url' = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
'userName' = '<yourUsername>',
'password' = '<yourPassword>',
'tableName' = '<yourTablename>'
);
INSERT INTO adb_sink
SELECT * FROM datagen_source;
Tabel dimensi
Contoh berikut melakukan join antara sumber datagen dengan tabel dimensi AnalyticDB for MySQL menggunakan temporal join. Hasilnya ditulis ke sink blackhole.
CREATE TEMPORARY TABLE datagen_source (
`a` INT,
`b` VARCHAR,
`c` STRING,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE adb_dim (
`a` INT,
`b` VARCHAR,
`c` VARCHAR
) WITH (
'connector' = 'adb3.0',
'url' = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
'userName' = '<yourUsername>',
'password' = '<yourPassword>',
'tableName' = '<yourTablename>'
);
CREATE TEMPORARY TABLE blackhole_sink (
`a` INT,
`b` VARCHAR
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN adb_dim FOR SYSTEM_TIME AS OF T.proctime AS H
ON T.a = H.a;