Le connecteur Tair (compatible avec Redis OSS) permet de lire des données depuis Tair en tant que table de dimension et d'écrire dans Tair en tant que table de destination dans les jobs de streaming Flink SQL.
Présentation du connecteur
| Élément | Détails |
|---|---|
| Types de table | Table de dimension, table de destination |
| Mode d'exécution | Streaming |
| Format de données | STRING |
| Métriques | Tables de dimension : aucune. Tables de destination : numBytesOut, numRecordsOutPerSecond, numBytesOutPerSecond, currentSendTime. Pour plus de détails, consultez la section Métriques de surveillance. |
| Type d'API | API SQL |
| Mises à jour et suppressions de données dans les tables de destination | Prises en charge |
Prérequis
Avant de commencer, vérifiez que vous disposez des éléments suivants :
Une instance Tair (compatible avec Redis OSS). Consultez la section Étape 1 : Créer une instance.
Une liste d'autorisation configurée pour l'instance. Consultez la section Étape 2 : Configurer les listes d'autorisation.
Limites
Sémantique de livraison : Le connecteur prend uniquement en charge la livraison au mieux (best-effort). La sémantique exactly-once n'est pas prise en charge. Assurez-vous que vos opérations d'écriture sont idempotentes.
Types de données des tables de dimension : Les tables de dimension ne peuvent lire que des données de type STRING et HASHMAP. Tous les champs doivent être de type STRING.
Clé primaire de la table de dimension : Chaque table de dimension doit comporter exactement une clé primaire. La clause ON d'une jointure avec une table de dimension doit utiliser des conditions d'égalité sur la clé primaire.
Problèmes connus
VVR 8.0.9 — Bug du cache Buffered Writer : Un problème lié au cache Buffered Writer existe dans VVR 8.0.9. Pour contourner ce problème, définissez sink.buffer-flush.max-rows sur 0 dans la clause WITH de la table de destination.
Syntaxe
CREATE TABLE redis_table (
col1 STRING,
col2 STRING,
PRIMARY KEY (col1) NOT ENFORCED -- Required.
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'mode' = 'STRING' -- Required for sink tables.
);
Options du connecteur
Options générales
Ces options s'appliquent aux tables de destination et aux tables de dimension.
| Option | Type de données | Obligatoire | Valeur par défaut | Description |
|---|---|---|---|---|
connector |
STRING | Oui | — | Définissez cette option sur redis. |
host |
STRING | Oui | — | Adresse IP utilisée pour se connecter à la base de données ApsaraDB for Redis. Utilisez l'endpoint interne chaque fois que possible. Les connexions Internet peuvent présenter une latence plus élevée ou des limitations de bande passante. |
port |
INT | Non | 6379 |
Numéro de port. |
password |
STRING | Non | (chaîne vide) | Mot de passe d'accès. |
dbNum |
INT | Non | 0 |
Numéro de séquence de la base de données. |
clusterMode |
BOOLEAN | Non | false |
Indique si la base de données est en mode cluster. |
hostAndPorts |
STRING | Non | — | Paires hôte et port au format "host1:port1,host2:port2". Obligatoire lorsque `clusterMode` est défini sur `true` et qu'une haute disponibilité (HA) est requise pour les connexions Jedis au cluster Redis auto-géré. Cette option est prioritaire sur `host` et `port`. Si `clusterMode` est défini sur `true` mais que la haute disponibilité n'est pas requise, configurez uniquement `host` et `port` pour spécifier un seul nœud. |
key-prefix |
STRING | Non | — | Préfixe ajouté à la valeur de la clé primaire lors de la lecture depuis une table de dimension ou de l'écriture dans une table de destination. Un délimiteur spécifié par `key-prefix-delimiter` sépare le préfixe et la valeur de la clé primaire. Nécessite VVR 8.0.7 ou version ultérieure. |
key-prefix-delimiter |
STRING | Non | — | Délimiteur entre le préfixe de clé et la valeur de la clé primaire. |
connection.pool.max-total |
INT | Non | 8 |
Nombre maximal de connexions pouvant être allouées par le pool de connexions. Nécessite VVR 8.0.9 ou version ultérieure. |
connection.pool.max-idle |
INT | Non | 8 |
Nombre maximal de connexions inactives dans le pool de connexions. |
connection.pool.min-idle |
INT | Non | 0 |
Nombre minimal de connexions inactives dans le pool de connexions. |
connection.pool.lifo |
Boolean | Non | true |
Indique si les connexions inactives sont allouées depuis le pool de connexions selon l'ordre LIFO (Last In, First Out). Valeurs valides :
Remarque : Cette option est prise en charge uniquement dans le moteur Realtime Compute VVR 11.8.0 et versions ultérieures. |
connect.timeout |
DURATION | Non | 3000ms |
Délai d'attente pour l'établissement de la connexion. |
socket.timeout |
DURATION | Non | 3000ms |
Délai d'attente pour la réception des données du serveur Redis. |
cacert.filepath |
STRING | Non | — | Chemin complet vers le certificat SSL/TLS. Le fichier doit être au format JKS. Si cette option n'est pas définie, le chiffrement SSL/TLS est désactivé. Pour activer le chiffrement, téléchargez le certificat CA et téléchargez-le en tant que dépendance supplémentaire — il est stocké dans le répertoire /flink/usrlib. Exemple : 'cacert.filepath' = '/flink/usrlib/ca.jks'. Nécessite VVR 11.1 ou version ultérieure. |
Options de la table de destination
| Option | Type de données | Obligatoire | Valeur par défaut | Description |
|---|---|---|---|---|
mode |
STRING | Oui | — | Structure de données Redis pour la table de destination. Cinq structures sont prises en charge : STRING, LIST, SET, HASHMAP et SORTEDSET. L'instruction DDL doit correspondre à la structure choisie. Consultez la section Structures de données pour les tables de destination. |
flattenHash |
BOOLEAN | Non | false |
Indique s'il faut écrire les données HASHMAP en mode multivaleur. Lorsque true est défini, déclarez plusieurs champs non primaires : la clé primaire correspond à la clé Redis, chaque nom de champ non primaire correspond à un champ Hash et chaque valeur de champ correspond à la valeur Hash. Lorsque false (mode monovaleur), déclarez exactement trois champs : la clé primaire correspond à la clé, le premier champ non primaire correspond au champ Hash et le second correspond à la valeur Hash. Prend effet uniquement lorsque mode est défini sur HASHMAP. Nécessite VVR 8.0.7 ou version ultérieure. |
ignoreDelete |
BOOLEAN | Non | false |
Indique s'il faut ignorer les messages de rétraction. Lorsque true est défini, les messages de rétraction sont ignorés. Lorsque false est défini, la clé et ses données sont supprimées lors de la réception d'un message de rétraction. |
expiration |
LONG | Non | 0 |
Durée de vie (TTL) des clés insérées, en millisecondes. La valeur 0 désactive la TTL. |
sink.buffer-flush.max-rows |
INT | Non | 200 |
Nombre maximal d'enregistrements (événements d'ajout, de modification et de suppression) conservés dans le tampon avant sa vidange. Nécessite VVR 8.0.9 ou version ultérieure pour clusterMode = false ; VVR 11.4.0 ou version ultérieure pour clusterMode = true. |
sink.buffer-flush.interval |
DURATION | Non | 1000ms |
Intervalle de vidange asynchrone du tampon. Nécessite VVR 8.0.9 ou version ultérieure pour clusterMode = false ; VVR 11.4.0 ou version ultérieure pour clusterMode = true. |
Options de la table de dimension
| Option | Type de données | Obligatoire | Valeur par défaut | Description |
|---|---|---|---|---|
mode |
STRING | Non | STRING |
Type de données à lire depuis la table de dimension. La valeur STRING lit les données STRING. La valeur HASHMAP lit les données Hash imbriquées (Key → Map\hashName. |
hashName |
STRING | Non | — | Clé Hash fixe utilisée lors de la lecture des données HASHMAP en mode monovaleur. Lorsqu'elle est définie, déclarez deux champs : la clé primaire correspond au champ Hash et la clé non primaire correspond à la valeur Hash. |
cache |
STRING | Non | None |
Politique de mise en cache. La valeur None désactive la mise en cache. La valeur LRU met en cache un sous-ensemble de données — en cas d'absence dans le cache, le connecteur interroge la table de dimension. La valeur ALL charge l'intégralité de la table de dimension dans le cache avant l'exécution du déploiement ; toutes les recherches ultérieures utilisent le cache, qui est rechargé après l'expiration des entrées. Consultez la section Remarques d'utilisation de l'option cache. |
cacheSize |
LONG | Non | 10000 |
Nombre maximal de lignes à mettre en cache. Requis lorsque cache est défini sur LRU. |
cacheTTLMs |
LONG | Non | — | Délai d'expiration du cache en millisecondes. Pour LRU, définit l'expiration par entrée (aucune expiration par défaut). Pour ALL, définit l'intervalle de rechargement du cache (aucun rechargement par défaut). Sans effet lorsque cache est défini sur None. |
cacheEmpty |
BOOLEAN | Non | true |
Indique s'il faut mettre en cache les résultats vides (sans correspondance). |
cacheReloadTimeBlackList |
STRING | Non | — | Périodes pendant lesquelles la politique de cache ALL ne procède pas au rechargement. Utile lors d'événements à fort trafic. Format : utilisez -> pour séparer les heures de début et de fin, et , pour séparer plusieurs périodes. Exemples : journée unique 2017-10-24 14:00 -> 2017-10-24 15:00 ; jours croisés 2017-11-10 23:30 -> 2017-11-11 08:00 ; récurrence quotidienne 12:00 -> 14:00, 22:00 -> 2:00 (nécessite VVR 11.1 ou version ultérieure). |
async |
BOOLEAN | Non | false |
Indique s'il faut activer la recherche asynchrone. Lorsque true est défini, les résultats sont renvoyés hors ordre. |
Remarques d'utilisation de l'option cache
ALLnécessite VVR 8.0.3 ou version ultérieure.Pour les versions VVR 8.0.3 à VVR 11.1 (exclusive),
cache = ALLlit HASHMAP uniquement en mode monovaleur. Dans la clause WITH de l'instruction DDL, définissezhashNamesur le nom de la clé ; déclarez Field comme clé primaire et Value comme colonne non primaire.À partir de VVR 11.1,
cache = ALLprend en charge le mode multivaleur HASHMAP. Spécifiez la clé Redis comme clé primaire et déclarez plusieurs colonnes non primaires pour chaque champ Hash. Définissezmode = HASHMAPdans la clause WITH.L'option
cachedoit être utilisée conjointement aveccacheSizeetcacheTTLMs.
Structures de données pour les tables de destination
Chaque structure de données Redis correspond à un schéma DDL spécifique et à une commande d'écriture.
| Structure de données | Schéma DDL | Commande d'écriture |
|---|---|---|
| STRING | Deux colonnes : key (STRING), value (STRING) | set key value |
| LIST | Deux colonnes : key (STRING), value (STRING) | lpush key value |
| SET | Deux colonnes : key (STRING), value (STRING) | sadd key value |
| HASHMAP (mode monovaleur, par défaut) | Trois colonnes : key (STRING), field (STRING), value (STRING) | hmset key field value |
HASHMAP (mode multivaleur, flattenHash = true) |
Plusieurs colonnes : key (STRING), puis une colonne par champ Hash — chaque nom de colonne correspond au nom du champ et sa valeur correspond à la valeur du champ | hmset key col1 value1 col2 value2 ... |
| SORTEDSET | Trois colonnes : key (STRING), score (DOUBLE), value (STRING) | zadd key score value |
L'optionignoreDeletecontrôle la gestion des messages de rétraction. Lorsqu'elle est définie surtrue, les opérations de suppression sont ignorées.
Mappages de types de données
| Portée | Type Tair (ApsaraDB for Redis) | Type Flink |
|---|---|---|
| Tous les types de table | STRING | STRING |
| Tables de destination uniquement | SCORE | DOUBLE |
Le type SCORE est utilisé avec les données SORTEDSET. Chaque valeur d'un ensemble trié nécessite un score DOUBLE, et les valeurs sont triées par ordre croissant de leurs scores.
Exemples
Exemples de tables de destination
Tous les exemples de tables de destination lisent depuis une source Kafka et écrivent dans une destination Tair.
Écriture de données STRING
Cet exemple utilise user_id comme clé Redis et login_time comme valeur Redis.
CREATE TEMPORARY TABLE kafka_source (
user_id STRING, -- User ID
login_time STRING -- Login time (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, -- Redis key
login_time STRING, -- Redis value
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;
Écriture de données HASHMAP en mode multivaleur
Cet exemple utilise order_id comme clé Redis et écrit product_name, quantity et amount en tant que champs Hash distincts.
CREATE TEMPORARY TABLE kafka_source (
order_id STRING, -- Order ID
product_name STRING, -- Product name
quantity STRING, -- Product quantity
amount STRING -- Order amount
) 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, -- Redis key
product_name STRING, -- Hash field: product_name
quantity STRING, -- Hash field: quantity
amount STRING, -- Hash field: 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;
Écriture de données HASHMAP en mode monovaleur
Cet exemple utilise order_id comme clé Redis, product_name comme champ Hash et quantity comme valeur Hash.
CREATE TEMPORARY TABLE kafka_source (
order_id STRING, -- Order ID
product_name STRING, -- Product name
quantity STRING -- Product quantity
) 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, -- Redis key
product_name STRING, -- Redis Hash field
quantity STRING, -- Redis Hash value
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;
Exemples de tables de dimension
Tous les exemples de tables de dimension recherchent les informations utilisateur dans une table de dimension Tair et les joignent à un flux Kafka.
Lecture de données STRING
Cet exemple utilise user_id comme clé Redis et récupère user_name en tant que valeur 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, -- Redis key
user_name STRING, -- Redis value
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;
Lecture de données HASHMAP en mode multivaleur
Cet exemple utilise user_id comme clé Redis et récupère plusieurs champs Hash — user_name, email et register_time — en une seule recherche.
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, -- Redis key
user_name STRING, -- Hash field: user_name
email STRING, -- Hash field: email
register_time STRING, -- Hash field: 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;
Lecture de données HASHMAP en mode monovaleur
Cet exemple utilise une clé Hash fixe (testkey) définie via hashName. La colonne user_id correspond au champ Hash et user_name correspond à la valeur 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, -- Hash field
user_name STRING, -- Hash value
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>',
'hashName' = 'testkey' -- Fixed Hash key
);
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;