Le connecteur Tair (Enterprise Edition) écrit les données de streaming Flink dans une instance Tair (Enterprise Edition). Tair est une base de données compatible avec Redis qui prend en charge le stockage hybride mémoire et disque, la haute disponibilité via un standby actif, ainsi qu'une architecture de cluster extensible adaptée aux charges de travail à haut débit et faible latence.
Fonctionnalités prises en charge
| Catégorie | Détails |
|---|---|
| Types pris en charge | Table de destination (Sink) |
| Mode d'exécution | Mode Stream |
| Format des données | STRING |
| Type d'API | SQL |
| Mise à jour ou suppression des données dans une table de destination | Oui |
| Métriques uniques | numBytesSend, numBytesSendPerSecond, numRecordsSend, numRecordsSendPerSecond, numRecordSendErrors, currentSendTime |
Pour plus de détails sur les métriques, consultez la rubrique Métriques.
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Une instance Tair (Enterprise Edition). Voir Étape 1 : Créer une instance
Une liste d'autorisation IP configurée pour l'instance. Voir Étape 2 : Configurer les listes d'autorisation
Limitations
Realtime Compute for Apache Flink nécessite Ververica Runtime (VVR) version 6.0.6 ou ultérieure pour utiliser le connecteur Tair (Enterprise Edition).
TairTs, TairCpc, TairRoaring, TairVector et TairGis nécessitent VVR version 8.0.1 ou ultérieure.
Le connecteur ne prend pas en charge la configuration de plusieurs hôtes.
Syntaxe
L'instruction DDL suivante crée une table de destination Tair. La clause PRIMARY KEY est obligatoire.
CREATE TABLE tair_table (
a STRING,
b STRING,
PRIMARY KEY (a) NOT ENFORCED
) WITH (
'connector' = 'tair',
'host' = '<yourHost>',
'mode' = '<dataStructure>'
);
Tair prend en charge toutes les structures de données Redis (STRING, LIST, SET, HASHMAP, SORTEDSET) ainsi que ses propres structures de données propriétaires. Pour des exemples de syntaxe concernant les structures de données compatibles avec Redis, consultez le connecteur Tair (compatible avec Redis OSS).
Options du connecteur
Paramètres obligatoires
| Paramètre | Type | Description |
|---|---|---|
connector |
String | Définissez la valeur sur tair. |
host |
String | Endpoint du serveur Tair. Utilisez un endpoint interne pour éviter la latence et les limitations de bande passante liées aux connexions via le réseau public. |
mode |
STRING | Structure de données cible pour l'écriture. Consultez la section Structures de données prises en charge pour connaître les valeurs valides. |
Paramètres de connexion
| Paramètre | Type | Par défaut | Description |
|---|---|---|---|
port |
INT | 6379 |
Numéro de port du serveur Tair. |
password |
String | Chaîne vide | Mot de passe de la base de données Tair. Laissez ce champ vide pour désactiver la vérification du mot de passe. |
dbNum |
INT | 0 |
Identifiant de la base de données de destination. |
clusterMode |
BOOLEAN | false |
Définissez la valeur sur true pour utiliser l'architecture en cluster. Définissez-la sur false pour le mode autonome. |
Paramètres de comportement d'écriture
| Paramètre | Type | Par défaut | Description |
|---|---|---|---|
ignoreDelete |
BOOLEAN | false |
Contrôle le comportement lors de la réception d'un message de rétractation. false : supprime les données insérées et leur clé. true : conserve les données insérées et leur clé. |
incrMode |
STRING | None |
Mode d'écriture de la destination. None : opération d'insertion. int : INCRBY avec un incrément fixe défini par incrValue. float : INCRBYFLOAT avec un incrément fixe défini par incrValue. dynamic_int : INCRBY avec l'incrément lu depuis la colonne nommée par incrValue. dynamic_float : INCRBYFLOAT avec l'incrément lu depuis la colonne nommée par incrValue. |
incrValue |
STRING | None | Lorsque incrMode est défini sur int ou float, il s'agit de la valeur d'incrément fixe. Lorsque incrMode est défini sur dynamic_int ou dynamic_float, il s'agit du nom de la colonne DDL contenant la valeur d'incrément. Non utilisé lorsque incrMode est défini sur None. |
Paramètres d'expiration des clés
| Paramètre | Type | Par défaut | Description |
|---|---|---|---|
expiration |
LONG | 0 |
Durée de vie relative (TTL) des clés insérées, en millisecondes. La valeur 0 désactive le TTL. |
expirationAt |
LONG | 0 |
Horodatage d'expiration absolue des clés insérées, en millisecondes. Ce paramètre n'est effectif que si expiration est défini sur 0. La valeur 0 désactive l'expiration absolue. |
Paramètres d'expiration au niveau des champs (TairHash et TairTs uniquement)
| Paramètre | Type | Par défaut | Description |
|---|---|---|---|
fieldExpireMode |
String | None |
Mode d'expiration des champs dans TairHash, ou des skeys dans TairTs. None : aucune expiration. millisecond : expiration relative définie sur la valeur de fieldExpireValue. unixtime : expiration absolue définie sur la valeur de fieldExpireValue. dynamic_millisecond : expiration relative lue depuis la colonne nommée par fieldExpireValue. dynamic_unixtime : expiration absolue lue depuis la colonne nommée par fieldExpireValue. Important
Les skeys TairTs doivent utiliser |
fieldExpireValue |
String | None | Lorsque fieldExpireMode est défini sur millisecond ou unixtime, il s'agit de la valeur d'expiration fixe. Lorsque fieldExpireMode est défini sur dynamic_millisecond ou dynamic_unixtime, il s'agit du nom de la colonne DDL contenant la valeur d'expiration. |
Mappages de types de données
| Type Flink | Type Tair |
|---|---|
| VARCHAR | STRING |
| DOUBLE | DOUBLE |
Structures de données prises en charge
Le paramètre mode accepte les valeurs suivantes. Chaque structure de données impose des exigences spécifiques concernant les colonnes DDL, qui varient selon le incrMode.
Structures de données compatibles avec Redis
Pour les formats DDL des structures STRING, LIST, SET, HASHMAP et SORTEDSET, consultez le connecteur Tair (compatible avec Redis OSS).
Structures de données propriétaires Tair
| Structure de données | **incrMode** |
Colonnes DDL | Commande d'écriture |
|---|---|---|---|
| TairString | None |
2 : key (STRING), value (STRING) | exset key value |
| TairString | int ou float |
1 : key (STRING) | exincrby/exincrbyfloat key incrValue |
| TairString | dynamic_int ou dynamic_float |
2 : key (STRING), incrValue (STRING) | exincrby/exincrbyfloat key incrValue |
| TairHash | None |
3 : key (STRING), field (STRING), value (STRING) | exhset key field value |
| TairHash | int ou float |
2 : key (STRING), field (STRING) | exhincrby/exincrbyfloat key field incrValue |
| TairHash | dynamic_int ou dynamic_float |
3 : key (STRING), field (STRING), incrValue (STRING) | exhincrby/exincrbyfloat key field incrValue |
| TairZset | None |
3–258 : key (STRING), member (STRING), score×N (DOUBLE, jusqu'à 256 dimensions) | exzadd key score member |
| TairZset | int ou float |
2 : key (STRING), member (STRING) | exzincyby key member incrValue |
| TairZset | dynamic_int ou dynamic_float |
3 : key (STRING), member (STRING), incrValue (STRING) | exzincyby key member incrValue |
| TairBloom | Doit être None |
2 : key (STRING), item (STRING) | BF.ADD key item |
| TairDoc | Doit être None |
3 : key (STRING), path (STRING), json (STRING) | JSON.SET key path json |
| TairSearch | None |
4 : index (STRING), doc_id (STRING), document (STRING, JSON), mapping (STRING) | TFT.ADDDOC index document docid |
| TairSearch | int ou float |
4 : index (STRING), doc_id (STRING), field (STRING), mapping (STRING) | TFT.INCRLONGDOCFIELD/TFT.INCRFLOATDOCFIELD index doc_id field increment |
| TairSearch | dynamic_int ou dynamic_float |
5 : index (STRING), doc_id (STRING), field (STRING), mapping (STRING), incrValue (STRING) | TFT.INCRLONGDOCFIELD/TFT.INCRFLOATDOCFIELD index doc_id field increment |
| TairCpc | Doit être None |
2 : key (STRING), item (STRING) | CPC.UPDATE key item |
| TairGis | Doit être None |
3 : key (STRING), polygon_name (STRING), polygon_wkt (STRING) | GIS.ADD area polygonName polygonWkt |
| TairRoaring | Doit être None |
3 : key (STRING), offset (BIGINT), value (BIGINT, 0 ou 1) | TR.SETBIT key offset value |
| TairVector | Doit être None |
6 : index_name (STRING), pk (STRING), vector_data (STRING), dims (INT), algorithm (STRING), distance_method (STRING) | TVS.HSET index_name key VECTOR vector_data |
| TairTs | None |
4 : pkey (STRING, groupe de chronologie), skey (STRING, chronologie unique), timestamp (STRING), value (STRING) | EXTS.S.RAW_MODIFY Pkey Skey timestamp value |
| TairTs | float |
3 : pkey (STRING), skey (STRING), timestamp (STRING) | EXTS.S.RAW_INCRBY Pkey Skey timestamp incrValue |
| TairTs | dynamic_float |
4 : pkey (STRING), skey (STRING), timestamp (STRING), incrValue (STRING) | EXTS.S.RAW_INCRBY Pkey Skey timestamp incrValue |
TairSearch et TairVector nécessitent la création préalable d'un index et d'un mappage avant l'insertion des données. Pour TairSearch, exécutezTFT.CREATEINDEX index mappings. Pour TairVector, exécutezTVS.CREATEINDEX index_name dims algorithm distance_method. TairBloom crée une clé avec une capacité par défaut de 100 éléments et un taux d'erreur de 0,01 lors de la première insertion. Pour le tri multidimensionnel avec TairZset, toutes les dimensions de score doivent utiliser le même format.
Exemples
Écrire des données dans une table de destination TairSearch
Cet exemple utilise index_name comme index TairSearch, doc_id comme identifiant du document, doc comme corps du document JSON et mapping comme définition du mappage de l'index.
CREATE TEMPORARY TABLE datagen_stream (
v STRING, -- document content
p STRING -- mapping definition
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
index_name STRING, -- TairSearch index name
doc_id STRING, -- document ID
doc STRING, -- document body (JSON)
mapping STRING, -- index mapping
PRIMARY KEY (index_name) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairsearch',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}'
);
INSERT INTO tair_output
SELECT
'index' AS index_name,
v AS doc_id,
p AS doc,
'{"mappings":{"_source":{"enabled":true},"properties":{"product_id":{"type":"keyword","ignore_above":128},"product_name":{"type":"text"}}}}' AS mapping
FROM datagen_stream;
Écrire des données avec un incrément dynamique (TairString avec incrMode=dynamic_float)
Cet exemple utilise key comme clé TairString et step comme valeur d'incrément par enregistrement, lue depuis la colonne step.
CREATE TEMPORARY TABLE datagen_stream (
v STRING, -- key
p STRING -- increment value per record
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
key STRING, -- TairString key
step STRING, -- increment value (column referenced by incrValue)
PRIMARY KEY (key) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairstring',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}',
'incrMode' = 'dynamic_float',
'incrValue' = 'step'
);
INSERT INTO tair_output
SELECT *
FROM datagen_stream;
Écrire des données avec un incrément fixe (TairString avec incrMode=float)
Cet exemple applique un incrément fixe de 11.11 à chaque clé lors de chaque écriture.
CREATE TEMPORARY TABLE datagen_stream (
v STRING,
p STRING
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
key STRING, -- TairString key
PRIMARY KEY (key) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairstring',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}',
'incrMode' = 'float',
'incrValue' = '11.11'
);
INSERT INTO tair_output
SELECT v
FROM datagen_stream;
Écrire des données avec une expiration au niveau du champ (TairHash avec fieldExpireMode=millisecond)
Cet exemple écrit dans une table de destination TairHash où chaque champ expire 1000 millisecondes après son écriture.
CREATE TEMPORARY TABLE datagen_stream (
v STRING,
p STRING,
s STRING
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
key STRING, -- TairHash key
field STRING, -- hash field
value STRING, -- field value
PRIMARY KEY (key) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairhash',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}',
'fieldExpireMode' = 'millisecond',
'fieldExpireValue' = '1000'
);
INSERT INTO tair_output
SELECT v, p, s
FROM datagen_stream;
Étapes suivantes
Connecteur Tair (compatible avec Redis OSS) — exemples de syntaxe pour les structures de données compatibles avec Redis
Métriques — liste complète des métriques du connecteur et leurs définitions