All Products
Search
Document Center

Realtime Compute for Apache Flink:Konektor Tair (Redis OSS-compatible)

Last Updated:Jul 31, 2026

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:

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:
  • true: LIFO. Kolam mengalokasikan koneksi yang paling baru dikembalikan terlebih dahulu.
  • false: FIFO. Koneksi idle yang paling lama tidak digunakan dialokasikan terlebih dahulu.
Untuk instans kluster Redis di balik proxy, mengatur ini ke false dapat menyeimbangkan beban di seluruh proxy secara lebih merata dan menghindari ketimpangan koneksi. Namun, FIFO cenderung membuat koneksi tetap aktif dan dapat meningkatkan jumlah total koneksi; evaluasi hal ini secara hati-hati dalam skenario dengan banyak koneksi.

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

  • ALL memerlukan VVR 8.0.3 atau versi lebih baru.

  • Untuk VVR 8.0.3 hingga VVR 11.1 (tidak termasuk), cache = ALL hanya membaca HASHMAP dalam mode nilai tunggal. Dalam klausa WITH DDL, atur hashName ke nama kunci; deklarasikan Field sebagai kunci primer dan Value sebagai kolom non-kunci primer.

  • Mulai VVR 11.1, cache = ALL mendukung mode multi-nilai HASHMAP. Tentukan kunci Redis sebagai kunci primer dan deklarasikan beberapa kolom non-kunci primer untuk setiap bidang Hash. Atur mode = HASHMAP dalam klausa WITH.

  • Opsi cache harus digunakan bersama dengan cacheSize dan cacheTTLMs.

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
Opsi ignoreDelete mengontrol cara penanganan pesan retraction. Ketika diatur ke true, 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;