Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur Tair (compatible avec Redis OSS)

Dernière mise à jour :Aug 09, 2026

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 :

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 :
  • true : LIFO. Le pool alloue en premier la connexion retournée le plus récemment.
  • false : FIFO. La connexion inactive la moins récemment utilisée est allouée en premier.
Pour une instance de cluster Redis située derrière un proxy, définir cette option sur false peut équilibrer la charge plus uniformément entre les proxies et éviter les déséquilibres de connexion. Toutefois, le mode FIFO tend à maintenir les connexions occupées et peut augmenter le nombre total de connexions ; évaluez soigneusement cette configuration dans les scénarios impliquant un grand nombre de connexions.

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\) : déclarez plusieurs colonnes non primaires, où la clé primaire correspond à la clé Redis, chaque nom de colonne non primaire correspond à un champ Hash et sa valeur correspond à la valeur du champ. Nécessite VVR 8.0.7 ou version ultérieure. Pour lire les données HASHMAP en mode monovaleur, définissez plutôt 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

  • ALL nécessite VVR 8.0.3 ou version ultérieure.

  • Pour les versions VVR 8.0.3 à VVR 11.1 (exclusive), cache = ALL lit HASHMAP uniquement en mode monovaleur. Dans la clause WITH de l'instruction DDL, définissez hashName sur le nom de la clé ; déclarez Field comme clé primaire et Value comme colonne non primaire.

  • À partir de VVR 11.1, cache = ALL prend 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éfinissez mode = HASHMAP dans la clause WITH.

  • L'option cache doit être utilisée conjointement avec cacheSize et cacheTTLMs.

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'option ignoreDelete contrôle la gestion des messages de rétraction. Lorsqu'elle est définie sur true , 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;