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
cacheSizelignes (les lignes les plus anciennes sont évincées en premier).L'entrée est restée dans le cache pendant plus de
cacheTTLMsmillisecondes.
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,
c1etc2n'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 |