Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Tair (Enterprise)

Dernière mise à jour :Aug 09, 2026

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 :

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 millisecond.

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écutez TFT.CREATEINDEX index mappings. Pour TairVector, exécutez TVS.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