Konektor Tair (Redis OSS-compatible) memungkinkan Anda membaca dari Tair sebagai tabel dimensi dan menulis ke Tair sebagai tabel sink dalam pekerjaan streaming Flink SQL.
Ikhtisar konektor
| Item | Detail |
|---|---|
| Jenis tabel | Tabel dimensi, tabel sink |
| Mode eksekusi | Streaming |
| Format data | STRING |
| Metrik | Tabel dimensi: tidak ada. Tabel sink: numBytesOut, numRecordsOutPerSecond, numBytesOutPerSecond, currentSendTime. Untuk detailnya, lihat Metrik pemantauan. |
| Jenis API | SQL API |
| Pembaruan dan penghapusan data di tabel sink | Didukung |
Prasyarat
Sebelum memulai, pastikan Anda telah:
-
Membuat instans Tair (Redis OSS-compatible). Lihat Langkah 1: Buat instans.
-
Mengonfigurasi daftar putih untuk instans tersebut. Lihat Langkah 2: Konfigurasi daftar putih.
Batasan
-
Semantik pengiriman: Konektor hanya mendukung pengiriman terbaik (best-effort delivery). Semantik exactly-once tidak didukung. Pastikan operasi tulis Anda bersifat idempoten.
-
Tipe data tabel dimensi: Tabel dimensi hanya dapat membaca data STRING dan HASHMAP. Semua bidang harus bertipe STRING.
-
Kunci primer tabel dimensi: Setiap tabel dimensi harus memiliki tepat satu kunci primer. Klausul ON pada JOIN tabel dimensi harus menggunakan kondisi kesamaan pada kunci primer.
Isu yang diketahui
VVR 8.0.9 — Bug cache Buffered Writer: Terdapat masalah cache Buffered Writer pada VVR 8.0.9. Sebagai solusi sementara, atur sink.buffer-flush.max-rows menjadi 0 dalam klausa WITH tabel sink.
Sintaks
CREATE TABLE redis_table (
col1 STRING,
col2 STRING,
PRIMARY KEY (col1) NOT ENFORCED -- Wajib.
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'mode' = 'STRING' -- Wajib untuk tabel sink.
);
Opsi konektor
Opsi umum
Opsi berikut berlaku untuk tabel sink maupun tabel dimensi.
| Opsi | Tipe data | Wajib | Bawaan | Deskripsi |
|---|---|---|---|---|
connector |
STRING | Ya | — | Atur ke redis. |
host |
STRING | Ya | — | Alamat IP yang digunakan untuk menghubungkan ke database ApsaraDB for Redis. Gunakan Titik akhir internal bila memungkinkan. Koneksi Internet mungkin mengalami latensi lebih tinggi atau batasan bandwidth. |
port |
INT | Tidak | 6379 |
Nomor port. |
password |
STRING | Tidak | (string kosong) | Password akses. |
dbNum |
INT | Tidak | 0 |
Nomor urut database. |
clusterMode |
BOOLEAN | Tidak | false |
Apakah database berada dalam mode kluster. |
hostAndPorts |
STRING | Tidak | — | Pasangan host dan port dalam format "host1:port1,host2:port2". Diperlukan ketika clusterMode bernilai true dan diperlukan ketersediaan tinggi (HA) untuk koneksi Jedis ke kluster Redis yang dikelola sendiri. Memiliki prioritas lebih tinggi daripada host dan port. Jika clusterMode bernilai true tetapi HA tidak diperlukan, Anda cukup mengonfigurasi host dan port untuk menentukan node tunggal. |
key-prefix |
STRING | Tidak | — | Awalan yang ditambahkan ke nilai kunci primer saat membaca dari tabel dimensi atau menulis ke tabel sink. Pemisah yang ditentukan oleh key-prefix-delimiter memisahkan awalan dan nilai kunci primer. Memerlukan VVR 8.0.7 atau versi lebih baru. |
key-prefix-delimiter |
STRING | Tidak | — | Pemisah antara awalan kunci dan nilai kunci primer. |
connection.pool.max-total |
INT | Tidak | 8 |
Jumlah maksimum koneksi yang dapat dialokasikan oleh kolam koneksi. Memerlukan VVR 8.0.9 atau versi lebih baru. |
connection.pool.max-idle |
INT | Tidak | 8 |
Jumlah maksimum koneksi idle dalam kolam koneksi. |
connection.pool.min-idle |
INT | Tidak | 0 |
Jumlah minimum koneksi idle dalam kolam koneksi. |
connection.pool.lifo |
Boolean | Tidak | true |
Apakah koneksi idle dialokasikan dari kolam koneksi dalam urutan LIFO. Nilai yang valid:
Catatan: Opsi ini hanya didukung pada engine Realtime Compute VVR 11.8.0 dan versi lebih baru. |
connect.timeout |
DURATION | Tidak | 3000ms |
Timeout untuk penyiapan koneksi. |
socket.timeout |
DURATION | Tidak | 3000ms |
Timeout untuk menerima data dari server Redis. |
cacert.filepath |
STRING | Tidak | — | Path lengkap ke sertifikat SSL/TLS. File harus dalam format JKS. Jika tidak diatur, enkripsi SSL/TLS dinonaktifkan. Untuk mengaktifkan enkripsi, unduh sertifikat CA dan unggah sebagai dependensi tambahan — file tersebut disimpan di direktori /flink/usrlib. Contoh: 'cacert.filepath' = '/flink/usrlib/ca.jks'. Memerlukan VVR 11.1 atau versi lebih baru. |
Opsi tabel sink
| Opsi | Tipe data | Wajib | Bawaan | Deskripsi |
|---|---|---|---|---|
mode |
STRING | Ya | — | Struktur data Redis untuk tabel sink. Lima struktur didukung: STRING, LIST, SET, HASHMAP, dan SORTEDSET. Pernyataan DDL harus sesuai dengan struktur yang dipilih. Lihat Struktur data untuk tabel sink. |
flattenHash |
BOOLEAN | Tidak | false |
Apakah akan menulis data HASHMAP dalam mode multi-nilai. Ketika true, deklarasikan beberapa bidang non-kunci primer: kunci primer dipetakan ke kunci Redis, setiap nama bidang non-kunci primer dipetakan ke bidang Hash, dan setiap nilai bidang dipetakan ke nilai Hash. Ketika false (mode nilai tunggal), deklarasikan tepat tiga bidang: kunci primer dipetakan ke kunci, bidang non-kunci primer pertama dipetakan ke bidang Hash, dan yang kedua dipetakan ke nilai Hash. Hanya berlaku ketika mode bernilai HASHMAP. Memerlukan VVR 8.0.7 atau versi lebih baru. |
ignoreDelete |
BOOLEAN | Tidak | false |
Apakah akan mengabaikan pesan retraction. Ketika true, pesan retraction dibuang. Ketika false, kunci dan datanya dihapus saat pesan retraction diterima. |
expiration |
LONG | Tidak | 0 |
Time-to-live (TTL) untuk kunci yang dimasukkan, dalam milidetik. 0 menonaktifkan TTL. |
sink.buffer-flush.max-rows |
INT | Tidak | 200 |
Jumlah maksimum catatan (event append, modifikasi, dan hapus) yang disimpan dalam buffer sebelum flush. Memerlukan VVR 8.0.9 atau versi lebih baru untuk clusterMode = false; VVR 11.4.0 atau versi lebih baru untuk clusterMode = true. |
sink.buffer-flush.interval |
DURATION | Tidak | 1000ms |
Interval flush buffer secara asinkron. Memerlukan VVR 8.0.9 atau versi lebih baru untuk clusterMode = false; VVR 11.4.0 atau versi lebih baru untuk clusterMode = true. |
Opsi tabel dimensi
| Opsi | Tipe data | Wajib | Bawaan | Deskripsi |
|---|---|---|---|---|
mode |
STRING | Tidak | STRING |
Tipe data yang dibaca dari tabel dimensi. STRING membaca data STRING. HASHMAP membaca data Hash bersarang (Kunci → Map\<Bidang, Nilai\>): deklarasikan beberapa kolom non-kunci primer, di mana kunci primer dipetakan ke kunci Redis, setiap nama kolom non-kunci primer dipetakan ke bidang Hash, dan nilainya dipetakan ke nilai bidang tersebut. Memerlukan VVR 8.0.7 atau versi lebih baru. Untuk membaca data HASHMAP dalam mode nilai tunggal, atur hashName sebagai gantinya. |
hashName |
STRING | Tidak | — | Kunci Hash tetap yang digunakan saat membaca data HASHMAP dalam mode nilai tunggal. Saat diatur, deklarasikan dua bidang: kunci primer dipetakan ke bidang Hash, dan bidang non-kunci primer dipetakan ke nilai Hash. |
cache |
STRING | Tidak | None |
Kebijakan cache. None menonaktifkan caching. LRU menyimpan sebagian data dalam cache — saat terjadi cache miss, konektor melakukan kueri ke tabel dimensi. ALL memuat seluruh tabel dimensi ke dalam cache sebelum penerapan dijalankan; semua pencarian berikutnya menggunakan cache, dan cache dimuat ulang setelah entri kedaluwarsa. Lihat Catatan penggunaan opsi cache. |
cacheSize |
LONG | Tidak | 10000 |
Jumlah maksimum baris yang dicache. Diperlukan ketika cache bernilai LRU. |
cacheTTLMs |
LONG | Tidak | — | Timeout cache dalam milidetik. Untuk LRU, mengatur waktu kedaluwarsa per entri (tidak ada kedaluwarsa secara bawaan). Untuk ALL, mengatur interval muat ulang cache (tidak ada muat ulang secara bawaan). Tidak berpengaruh ketika cache bernilai None. |
cacheEmpty |
BOOLEAN | Tidak | true |
Apakah akan menyimpan hasil kosong (tidak cocok) dalam cache. |
cacheReloadTimeBlackList |
STRING | Tidak | — | Periode waktu ketika kebijakan cache ALL tidak melakukan muat ulang. Berguna selama acara dengan trafik tinggi. Format: gunakan -> untuk memisahkan waktu mulai dan selesai, serta , untuk memisahkan beberapa periode. Contoh: satu hari 2017-10-24 14:00 -> 2017-10-24 15:00; lintas hari 2017-11-10 23:30 -> 2017-11-11 08:00; berulang harian 12:00 -> 14:00, 22:00 -> 2:00 (memerlukan VVR 11.1 atau versi lebih baru). |
async |
BOOLEAN | Tidak | false |
Apakah akan mengaktifkan lookup asinkron. Ketika true, hasil dikembalikan di luar urutan. |
Catatan penggunaan opsi cache
-
ALLmemerlukan VVR 8.0.3 atau versi lebih baru. -
Untuk VVR 8.0.3 hingga VVR 11.1 (tidak termasuk),
cache = ALLhanya membaca HASHMAP dalam mode nilai tunggal. Dalam klausa WITH DDL, aturhashNameke nama kunci; deklarasikan Field sebagai kunci primer dan Value sebagai kolom non-kunci primer. -
Mulai VVR 11.1,
cache = ALLmendukung mode multi-nilai HASHMAP. Tentukan kunci Redis sebagai kunci primer dan deklarasikan beberapa kolom non-kunci primer untuk setiap bidang Hash. Aturmode = HASHMAPdalam klausa WITH. -
Opsi
cacheharus digunakan bersama dengancacheSizedancacheTTLMs.
Struktur data untuk tabel sink
Setiap struktur data Redis berkorespondensi dengan skema DDL dan perintah tulis tertentu.
| Struktur data | Skema DDL | Write Command |
|---|---|---|
| STRING | Dua kolom: key (STRING), value (STRING) | set key value |
| LIST | Dua kolom: key (STRING), value (STRING) | lpush key value |
| SET | Dua kolom: key (STRING), value (STRING) | sadd key value |
| HASHMAP (mode nilai tunggal, bawaan) | Tiga kolom: key (STRING), field (STRING), value (STRING) | hmset key field value |
HASHMAP (mode multi-nilai, flattenHash = true) |
Beberapa kolom: key (STRING), lalu satu kolom per bidang Hash — setiap nama kolom adalah nama bidang dan nilainya adalah nilai bidang tersebut | hmset key col1 value1 col2 value2 ... |
| SORTEDSET | Tiga kolom: key (STRING), score (DOUBLE), value (STRING) | zadd key score value |
OpsiignoreDeletemengontrol cara penanganan pesan retraction. Ketika diatur ketrue, operasi hapus dilewati.
Pemetaan tipe data
| Ruang lingkup | Tipe Tair (ApsaraDB for Redis) | Tipe Flink |
|---|---|---|
| Semua jenis tabel | STRING | STRING |
| Hanya tabel sink | SCORE | DOUBLE |
Tipe SCORE digunakan dengan data SORTEDSET. Setiap nilai dalam sorted set memerlukan skor DOUBLE, dan nilai-nilai tersebut diurutkan secara ascending berdasarkan skornya.
Contoh
Contoh tabel sink
Semua contoh tabel sink membaca dari sumber Kafka dan menulis ke sink Tair.
Menulis data STRING
Contoh ini menggunakan user_id sebagai kunci Redis dan login_time sebagai nilai Redis.
CREATE TEMPORARY TABLE kafka_source (
user_id STRING, -- ID Pengguna
login_time STRING -- Waktu login (Unix timestamp)
) WITH (
'connector' = 'kafka',
'topic' = 'user_logins',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_sink (
user_id STRING, -- Kunci Redis
login_time STRING, -- Nilai Redis
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'mode' = 'STRING',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>'
);
INSERT INTO redis_sink
SELECT * FROM kafka_source;
Menulis data HASHMAP dalam mode multi-nilai
Contoh ini menggunakan order_id sebagai kunci Redis dan menulis product_name, quantity, dan amount sebagai bidang Hash terpisah.
CREATE TEMPORARY TABLE kafka_source (
order_id STRING, -- ID Pesanan
product_name STRING, -- Nama produk
quantity STRING, -- Jumlah produk
amount STRING -- Jumlah pesanan
) WITH (
'connector' = 'kafka',
'topic' = 'orders_topic',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_sink (
order_id STRING, -- Kunci Redis
product_name STRING, -- Bidang Hash: product_name
quantity STRING, -- Bidang Hash: quantity
amount STRING, -- Bidang Hash: amount
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'mode' = 'HASHMAP',
'flattenHash' = 'true',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>'
);
INSERT INTO redis_sink
SELECT * FROM kafka_source;
Menulis data HASHMAP dalam mode nilai tunggal
Contoh ini menggunakan order_id sebagai kunci Redis, product_name sebagai bidang Hash, dan quantity sebagai nilai Hash.
CREATE TEMPORARY TABLE kafka_source (
order_id STRING, -- ID Pesanan
product_name STRING, -- Nama produk
quantity STRING -- Jumlah produk
) WITH (
'connector' = 'kafka',
'topic' = 'orders_topic',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_sink (
order_id STRING, -- Kunci Redis
product_name STRING, -- Bidang Hash Redis
quantity STRING, -- Nilai Hash Redis
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'mode' = 'HASHMAP',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>'
);
INSERT INTO redis_sink
SELECT * FROM kafka_source;
Contoh tabel dimensi
Semua contoh tabel dimensi melakukan lookup informasi pengguna dari tabel dimensi Tair dan menggabungkannya dengan aliran Kafka.
Membaca data STRING
Contoh ini menggunakan user_id sebagai kunci Redis dan mengambil user_name sebagai nilai Redis.
CREATE TEMPORARY TABLE kafka_source (
user_id STRING,
proctime AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'user_clicks',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_dim (
user_id STRING, -- Kunci Redis
user_name STRING, -- Nilai Redis
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>',
'mode' = 'STRING'
);
CREATE TEMPORARY TABLE blackhole_sink (
user_id STRING,
redis_user_id STRING,
user_name STRING
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT
t1.user_id,
t2.user_id,
t2.user_name
FROM kafka_source AS t1
JOIN redis_dim FOR SYSTEM_TIME AS OF t1.proctime AS t2
ON t1.user_id = t2.user_id;
Membaca data HASHMAP dalam mode multi-nilai
Contoh ini menggunakan user_id sebagai kunci Redis dan mengambil beberapa bidang Hash — user_name, email, dan register_time — dalam satu kali lookup.
CREATE TEMPORARY TABLE kafka_source (
user_id STRING,
click_time TIMESTAMP(3),
proctime AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'user_clicks',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_dim (
user_id STRING, -- Kunci Redis
user_name STRING, -- Bidang Hash: user_name
email STRING, -- Bidang Hash: email
register_time STRING, -- Bidang Hash: register_time
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>',
'mode' = 'HASHMAP'
);
CREATE TEMPORARY TABLE blackhole_sink (
user_id STRING,
user_name STRING,
email STRING,
register_time STRING,
click_time TIMESTAMP(3)
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT
t1.user_id,
t2.user_name,
t2.email,
t2.register_time,
t1.click_time
FROM kafka_source AS t1
JOIN redis_dim FOR SYSTEM_TIME AS OF t1.proctime AS t2
ON t1.user_id = t2.user_id;
Membaca data HASHMAP dalam mode nilai tunggal
Contoh ini menggunakan kunci Hash tetap (testkey) yang diatur melalui hashName. Kolom user_id dipetakan ke bidang Hash, dan user_name dipetakan ke nilai Hash.
CREATE TEMPORARY TABLE kafka_source (
user_id STRING,
proctime AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'user_clicks',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_dim (
user_id STRING, -- Bidang Hash
user_name STRING, -- Nilai Hash
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>',
'hashName' = 'testkey' -- Kunci Hash tetap
);
CREATE TEMPORARY TABLE blackhole_sink (
user_id STRING,
redis_user_id STRING,
user_name STRING
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT
t1.user_id,
t2.user_id,
t2.user_name
FROM kafka_source AS t1
JOIN redis_dim FOR SYSTEM_TIME AS OF t1.proctime AS t2
ON t1.user_id = t2.user_id;