Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur Lindorm

Dernière mise à jour :Aug 09, 2026

Sink : Streaming Source de recherche : Mode synchrone/asynchrone

Le connecteur Lindorm permet aux jobs Flink en streaming d'écrire des données et d'effectuer des recherches dans les tables larges Lindorm via l'API SQL.

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • Un moteur de table large Lindorm et une table Lindorm. Pour plus d'informations, consultez la rubrique Créer une instance.

  • Une connectivité réseau entre le cluster Lindorm et l'espace de travail Flink (par exemple, les deux doivent se trouver dans le même cloud privé virtuel (VPC)).

Les tables HBase Lindorm ne sont pas prises en charge. Seule LindormTable est compatible.

Démarrage rapide

L'exemple suivant génère 10 lignes de données, recherche les lignes correspondantes dans une table de dimension Lindorm, puis écrit le résultat dans une table de destination Lindorm.

-- Source: generate 10 rows with sequential IDs 0-9
CREATE TEMPORARY TABLE example_source (
  id INT,
  proc_time AS PROCTIME()
) WITH (
  'connector' = 'datagen',
  'number-of-rows' = '10',
  'fields.id.kind' = 'sequence',
  'fields.id.start' = '0',
  'fields.id.end' = '9'
);

-- Dimension table: look up user details from Lindorm
CREATE TEMPORARY TABLE lindorm_hbase_dim (
  `id`    INT,
  `name`  VARCHAR,
  `birth` VARCHAR,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector'   = 'lindorm',
  'tablename'   = '${lindorm_dim_table}',
  'seedserver'  = '${lindorm_seed_server}',
  'namespace'   = 'default',
  'username'    = '${lindorm_username}',
  'password'    = '${lindorm_password}'
);

-- Sink table: write enriched records to Lindorm
CREATE TEMPORARY TABLE lindorm_hbase_sink (
  `id`    INT,
  `name`  VARCHAR,
  `birth` VARCHAR,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector'   = 'lindorm',
  'tablename'   = '${lindorm_sink_table}',
  'seedserver'  = '${lindorm_seed_server}',
  'namespace'   = 'default',
  'username'    = '${lindorm_username}',
  'password'    = '${lindorm_password}'
);

-- Temporal join: enrich source data with the dimension table and write to sink
INSERT INTO lindorm_hbase_sink
SELECT
  s.id,
  d.name,
  d.birth
FROM example_source AS s
JOIN lindorm_hbase_dim AS d FOR SYSTEM_TIME AS OF s.proc_time
  ON s.id = d.id;

Remplacez les espaces réservés avant l'exécution :

Espace réservé Description
${lindorm_dim_table} Nom de la table de dimension Lindorm
${lindorm_sink_table} Nom de la table de destination Lindorm
${lindorm_seed_server} Endpoint du serveur Lindorm au format host:port
${lindorm_username} Nom d'utilisateur Lindorm
${lindorm_password} Mot de passe Lindorm

Syntaxe

CREATE TABLE white_list (
  id     VARCHAR,
  name   VARCHAR,
  age    INT,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector'    = 'lindorm',
  'seedserver'   = '<host:port>',
  'namespace'    = '<yourNamespace>',
  'username'     = '<yourUsername>',
  'password'     = '<yourPassword>',
  'tableName'    = '<yourTableName>',
  'columnFamily' = '<yourColumnFamily>'
);

Options du connecteur

Général

Option Type Obligatoire Par défaut Description
connector String Oui Doit être défini sur lindorm.
seedserver String Oui Endpoint du serveur Lindorm au format host:port. Realtime Compute for Apache Flink utilise l'API ApsaraDB for HBase pour Java afin d'établir la connexion. Pour plus de détails, consultez la rubrique Utiliser Flink pour se connecter et utiliser LindormTable.
namespace String Oui Espace de noms de la base de données Lindorm.
username String Oui Nom d'utilisateur de la base de données Lindorm.
password String Oui Mot de passe de la base de données Lindorm.
tableName String Oui Nom de la table Lindorm.
columnFamily String Oui Nom de la famille de colonnes. Si aucune famille de colonnes n'a été spécifiée lors de la création de la table, saisissez f.
retryIntervalMs Integer Non 1000 Intervalle de nouvelle tentative pour les opérations de lecture ayant échoué, en millisecondes.
maxRetryTimes Integer Non 5 Nombre maximal de nouvelles tentatives pour les opérations de lecture ou d'écriture.

Spécifique au sink

Option Type Obligatoire Par défaut Description
bufferSize Integer Non 500 Nombre d'enregistrements à mettre en mémoire tampon avant l'écriture dans Lindorm.
flushIntervalMs Integer Non 2000 Temps maximal entre les vidages lorsque la mémoire tampon n'est pas pleine, en millisecondes.
ignoreDelete Boolean Non false Lorsque la valeur est true, les opérations de suppression sont ignorées.
dynamicColumnSink Boolean Non false Lorsque la valeur est true, active la fonctionnalité de table dynamique. Consultez la section Table dynamique.
excludeUpdateColumns String Non Liste séparée par des virgules des colonnes à exclure des mises à jour. Par exemple, a,b,c ignore les mises à jour des colonnes a, b et c. Nécessite VVR 8.0.9 ou version ultérieure.

Spécifique à la table de dimension

Option Type Obligatoire Par défaut Description
partitionedJoin Boolean Non false Lorsque la valeur est true, utilise la JoinKey pour le partitionnement afin d'améliorer le taux de succès du cache.
shuffleEmptyKey Boolean Non false Lorsque la valeur est true, achemine aléatoirement les clés vides en amont vers les nœuds en aval. Lorsque la valeur est false, elles sont acheminées vers le thread parallèle 0.
cache String Non None Politique de mise en cache. Valeurs valides : None (aucune mise en cache) et LRU (mise en cache des lignes récemment consultées).
cacheSize Integer Non 1000 Nombre maximal de lignes à mettre en cache. S'applique lorsque cache est défini sur LRU.
cacheTTLMs Integer Non Durée d'expiration de l'entrée de cache, en millisecondes. S'applique lorsque cache est défini sur LRU. Par défaut, les entrées n'expirent pas.
cacheEmpty Boolean Non true Lorsque la valeur est true, met en cache les résultats de recherche qui n'ont retourné aucune ligne.
async Boolean Non false Lorsque la valeur est true, active le mode de recherche asynchrone.
asyncLindormRpcTimeoutMs Integer Non 300000 Délai d'expiration RPC pour les recherches asynchrones, en millisecondes.

Cache de recherche

Par défaut, le connecteur effectue chaque recherche directement depuis Lindorm (cache = None). Activez la mise en cache LRU pour réduire la pression de lecture sur Lindorm lors de jointures à haut débit.

Avec cache = LRU, le connecteur stocke les lignes récemment consultées en mémoire. Lorsque la clé de jointure correspond à une ligne mise en cache, le connecteur renvoie la valeur en cache sans interroger Lindorm. Une entrée mise en cache est évincée dans les cas suivants :

  • Le cache atteint cacheSize lignes (les lignes les plus anciennes sont évincées en premier).

  • L'entrée est restée dans le cache pendant plus de cacheTTLMs millisecondes.

Il existe un compromis : un TTL plus long ou un cache plus grand réduit le trafic de lecture Lindorm, mais augmente le risque de servir des données obsolètes. Ajustez ces deux valeurs en fonction de vos exigences de débit et de la rapidité avec laquelle les données sous-jacentes changent.

Le connecteur Lindorm prend en charge les jointures de recherche un-à-plusieurs. Soyez particulièrement attentif aux stratégies de mise en cache et au débit lorsque les lignes de la table de dimension peuvent correspondre à plusieurs événements en amont.

Écritures idempotentes

Lorsqu'une table Lindorm possède une clé primaire, toutes les écritures utilisent la sémantique upsert : chaque enregistrement entrant insère soit une nouvelle ligne, soit met à jour la ligne existante correspondante. Cela rend les écritures idempotentes.

Les écritures idempotentes sont cruciales pour la tolérance aux pannes. Si un job Flink redémarre à partir d'un point de contrôle, il rejoue les messages depuis le dernier point de contrôle réussi. Comme les upserts Lindorm sont idempotents, les enregistrements rejoués produisent le même résultat que les écritures originales, sans doublons ni violations de contrainte.

Définissez une clé primaire dans votre DDL pour tirer parti des écritures idempotentes.

Table dynamique

Utilisez la fonctionnalité de table dynamique lorsque votre schéma évolue au moment de l'exécution : les colonnes sont créées dynamiquement en fonction des valeurs présentes dans les données, plutôt que d'être figées lors de la définition DDL. Un cas d'utilisation typique est le suivi des métriques horaires par jour, où les heures constituent les noms de colonnes et les jours les clés primaires :

Clé primaire 00:00 01:00
2025-06-01 45 32
2025-06-02 76 34

Règles DDL pour les tables dynamiques :

  • Les N premières colonnes constituent la clé primaire.

  • Les deux dernières colonnes doivent être de type VARCHAR.

  • L'avant-dernière colonne (c1) contient le nom de la colonne dynamique.

  • La dernière colonne (c2) contient la valeur de cette colonne.

  • Aucune colonne autre que la clé primaire, c1 et c2 n'est autorisée.

CREATE TABLE lindorm_dynamic_output (
  pk1 VARCHAR,
  pk2 VARCHAR,
  pk3 VARCHAR,
  c1  VARCHAR,  -- column name written to Lindorm
  c2  VARCHAR,  -- column value written to Lindorm
  PRIMARY KEY (pk1, pk2, pk3) NOT ENFORCED
) WITH (
  'connector'         = 'lindorm',
  'seedserver'        = '<host:port>',
  'namespace'         = '<yourNamespace>',
  'username'          = '<yourUsername>',
  'password'          = '<yourPassword>',
  'tableName'         = '<yourTableName>',
  'columnFamily'      = '<yourColumnFamily>',
  'dynamicColumnSink' = 'true'
);

À chaque écriture d'enregistrement, le connecteur ajoute ou met à jour une colonne dans la ligne Lindorm identifiée par <pk1, pk2, pk3>. Les autres colonnes de cette ligne restent inchangées.

Mappages de types de données

Toutes les données Lindorm sont stockées au format binaire. Le tableau suivant indique comment le convertisseur opère la conversion entre les types SQL Flink et les représentations binaires Lindorm.

Type SQL Flink Écriture dans Lindorm Lecture depuis Lindorm
CHAR / VARCHAR StringData::toBytes StringData::fromBytes
BOOLEAN Bytes::toBytes(boolean) Bytes::toBigDecimal
BINARY / VARBINARY Octets directs Octets directs
DECIMAL Bytes::toBytes(BigDecimal) Bytes::toBigDecimal
TINYINT Premier octet de byte[] bytes[0]
SMALLINT Bytes::toBytes(short) Bytes::toShort
INT Bytes::toBytes(int) Bytes::toInt
BIGINT Bytes::toBytes(long) Bytes::toLong
FLOAT Bytes::toBytes(float) Bytes::toFloat
DOUBLE Bytes::toBytes(double) Bytes::toDouble
DATE Bytes::toBytes(int) avec le nombre de jours écoulés depuis le 01/01/1970 Bytes::toInt → nombre de jours écoulés depuis le 01/01/1970
TIME Bytes::toBytes(int) avec le nombre de millisecondes écoulées depuis 00:00:00 Bytes::toInt → nombre de millisecondes écoulées depuis 00:00:00
TIMESTAMP Bytes::toBytes(long) avec le nombre de millisecondes écoulées depuis le 01/01/1970 00:00:00 Bytes::toLong → nombre de millisecondes écoulées depuis le 01/01/1970 00:00:00

Les méthodes Bytes se trouvent dans la classe com.alibaba.lindorm.client.core.utils.Bytes. Les méthodes StringData (lignes CHAR/VARCHAR) se trouvent dans la classe org.apache.flink.table.data.StringData.

Métriques

Les métriques de sink suivantes sont disponibles. Pour plus de détails, consultez la rubrique Métriques.

Métrique Description
numBytesOut Total des octets écrits dans le sink
numBytesOutPerSecond Octets écrits par seconde
numRecordsOut Total des enregistrements écrits dans le sink
numRecordsOutPerSecond Enregistrements écrits par seconde

FAQ

Erreurs de connexion Lindorm et solutions