Cette rubrique explique comment utiliser le connecteur ApsaraDB for HBase.
Informations générales
ApsaraDB for HBase est un service NoSQL intelligent hébergé dans le cloud, économique et compatible avec la version open source de HBase. Il offre une grande extensibilité. ApsaraDB for HBase présente plusieurs avantages, notamment des coûts de stockage réduits, un débit élevé, une mise à l'échelle simplifiée et un traitement intelligent des données. Ce service prend en charge les activités principales d'Alibaba, telles que les recommandations Taobao, la gestion des risques pour Ant Credit Pay, la publicité, les tableaux de bord de données, le suivi logistique de Cainiao, les historiques de transactions Alipay et les messages de l'application mobile Taobao. En tant que service entièrement géré, ApsaraDB for HBase fournit des fonctionnalités de niveau entreprise, dont le traitement de pétaoctets de données, la prise en charge d'une forte concurrence, une mise à l'échelle rapide en quelques secondes, une latence de réponse inférieure à la milliseconde, une haute disponibilité entre centres de données et une distribution mondiale.
Le tableau suivant décrit les fonctionnalités prises en charge par le connecteur ApsaraDB for HBase.
Élément | Description |
Type de table | Table de dimension et table de destination (sink) |
Mode d'exécution | Mode streaming |
Format des données | S.O. |
Métrique | |
Type d'API | API SQL |
Mise à jour ou suppression des données dans une table de destination | Prise en charge |
Prérequis
Un cluster ApsaraDB for HBase doit être acheté et une table ApsaraDB for HBase doit être créée. Pour savoir comment acheter un cluster ApsaraDB for HBase, consultez Acheter un cluster.
Une liste d'autorisation doit être configurée pour le cluster ApsaraDB for HBase. Pour plus d'informations, consultez Configurer une liste d'autorisation.
Notes d'utilisation
Avant d'utiliser le connecteur ApsaraDB for HBase, vérifiez le type de votre instance de base de données et assurez-vous que le connecteur sélectionné correspond bien à ce type. Une utilisation incorrecte du connecteur peut entraîner des problèmes inattendus.
Le connecteur ApsaraDB for HBase décrit dans cette rubrique est destiné aux instances ApsaraDB for HBase.
Les instances Lindorm sont compatibles avec Apache HBase. Utilisez le connecteur Lindorm pour les instances Lindorm. Pour plus d'informations, consultez Lindorm.
Si vous utilisez le connecteur ApsaraDB for HBase pour connecter Realtime Compute for Apache Flink à une base de données HBase open source, l'intégrité des données ne peut pas être garantie.
Syntaxe
CREATE TABLE hbase_table(
rowkey INT,
family1 ROW<q1 INT>,
family2 ROW<q2 STRING, q3 BIGINT>,
family3 ROW<q4 DOUBLE, q5 BOOLEAN, q6 STRING>
) WITH (
'connector'='cloudhbase',
'table-name'='<yourTableName>',
'zookeeper.quorum'='<yourZookeeperQuorum>'
);
Les familles de colonnes d'une table ApsaraDB for HBase doivent être déclarées avec le type ROW. Chaque nom de famille de colonnes devient le nom d'un champ de la ligne. Dans la syntaxe DDL, les familles de colonnes suivantes sont déclarées : family1, family2 et family3.
Chaque colonne d'une famille de colonnes correspond à un champ d'une ligne. Le nom de la colonne devient le nom du champ. Dans la syntaxe DDL, les colonnes q2 et q3 sont déclarées dans la famille de colonnes family2.
Outre les champs de type ROW, une seule colonne de type atomique (tel que STRING ou BIGINT) peut exister dans une table ApsaraDB for HBase. Ce champ de type atomique est considéré comme la clé de ligne (row key) de la table, comme illustré par
rowkeydans l'instruction DDL.La clé de ligne d'une table ApsaraDB for HBase doit être définie comme clé primaire de la table de destination. Si aucune clé primaire n'est définie, la clé de ligne sert de clé primaire.
Il suffit de déclarer les familles de colonnes et les colonnes nécessaires de la table ApsaraDB for HBase dans la table de destination.
Options du connecteur
-
Général
Option
Description
Type de données
Obligatoire
Valeur par défaut
Remarques
connector
Le type de la table.
String
Oui
Aucune valeur par défaut
Définissez la valeur sur
cloudhbase.table-name
Le nom de la table ApsaraDB for HBase.
String
Oui
Aucune valeur par défaut
Sans objet.
zookeeper.znode.quorum
L'URL utilisée pour accéder au service ZooKeeper d'ApsaraDB for HBase.
String
Oui
Aucune valeur par défaut
Sans objet.
zookeeper.znode.parent
Le répertoire racine d'ApsaraDB for HBase dans le service ZooKeeper.
String
Non
/hbaseCe paramètre prend effet uniquement dans l'édition Standard d'ApsaraDB for HBase.
userName
Le nom d'utilisateur utilisé pour accéder à la base de données.
String
Non
Aucune valeur par défaut
Ce paramètre prend effet uniquement dans l'édition Performance-enhanced d'ApsaraDB for HBase.
password
Le mot de passe utilisé pour accéder à la base de données.
String
Non
Aucune valeur par défaut
Ce paramètre prend effet uniquement dans l'édition Performance-enhanced d'ApsaraDB for HBase.
haclient.cluster.id
L'ID du cluster ApsaraDB for HBase en mode haute disponibilité (HA).
String
Non
Aucune valeur par défaut
Ce paramètre est requis uniquement lorsque vous accédez aux clusters de reprise après sinistre interzones. Il prend effet uniquement dans l'édition Performance-enhanced d'ApsaraDB for HBase.
retires.number
Le nombre de tentatives autorisées pour que le client ApsaraDB for HBase se connecte à la base de données ApsaraDB for HBase.
Integer
Non
31
Sans objet.
null-string-literal
Si le type de données d'un champ d'ApsaraDB for HBase est STRING et que les données du champ dans Realtime Compute for Apache Flink sont nulles,
null-string-literalest attribué au champ d'ApsaraDB for HBase et écrit dans la base de données ApsaraDB for HBase.String
Non
null
Sans objet.
-
Spécifique au Sink
Option
Description
Data type
Required
Default value
Remarks
sink.buffer-flush.max-size
Taille des données mises en cache en mémoire (en octets) avant l'écriture dans la base de données ApsaraDB for HBase. Une valeur plus élevée améliore les performances d'écriture, mais augmente la latence et la consommation de mémoire.
String
No
2MB
Unités : B, KB, MB ou GB (non sensibles à la casse). Si ce paramètre est défini sur 0, aucune donnée n'est mise en cache.
sink.buffer-flush.max-rows
Nombre d'enregistrements mis en cache en mémoire avant l'écriture dans la base de données ApsaraDB for HBase. Une valeur plus élevée améliore les performances d'écriture, mais augmente la latence et la consommation de mémoire.
Integer
No
1000
Si ce paramètre est défini sur 0, aucune donnée n'est mise en cache.
sink.buffer-flush.interval
Intervalle d'écriture des données mises en cache dans la base de données ApsaraDB for HBase. Ce paramètre contrôle la latence d'écriture.
Duration
No
1s
Unités : ms, s, min, h ou d. Si ce paramètre est défini sur 0, l'écriture périodique des données est désactivée.
dynamic.table
Indique s'il faut utiliser une table ApsaraDB for HBase prenant en charge les colonnes dynamiques.
Boolean
No
false
Valeurs valides :
true
false
sink.ignore-delete
Indique s'il faut ignorer les messages de rétractation.
Boolean
No
false
Si un flux contient des événements
DELETEouUPDATE_BEFORE, et que plusieurs tâches sink mettent à jour simultanément différents champs d'une table, des incohérences de données peuvent survenir.Par exemple, après la suppression d'un enregistrement, une autre tâche met à jour certains champs. Les champs non mis à jour deviennent alors null ou sont remplis avec des valeurs par défaut, ce qui entraîne des erreurs de données.
Pour éviter ce problème, définissez
sink.ignore-deletesurtrueafin d'ignorer les événementsDELETEetUPDATE_BEFOREen amont.RemarqueUPDATE_BEFOREfait partie du mécanisme de rétractation de Flink et sert à rétracter l'ancienne valeur lors d'une opération de mise à jour.Si
ignoreDeleteest défini surtrue, tous les événementsDELETEetUPDATE_BEFOREsont ignorés. Seuls les enregistrementsINSERTetUPDATE_AFTERsont traités.
sink.sync-write
Indique s'il faut écrire les données dans ApsaraDB for HBase en mode synchrone.
Boolean
No
true
Valeurs valides :
true : Les données sont écrites en mode synchrone. Dans ce mode, les données sont écrites séquentiellement, mais les performances d'écriture sont réduites.
false : Les données sont écrites en mode asynchrone. Dans ce mode, les données peuvent ne pas être écrites séquentiellement, mais les performances d'écriture sont améliorées.
sink.buffer-flush.batch-rows
Nombre d'enregistrements mis en cache en mémoire lors de l'écriture des données dans ApsaraDB for HBase en mode synchrone. Une valeur plus élevée améliore les performances d'écriture, mais augmente la latence et l'utilisation de la mémoire.
Integer
No
100
Ce paramètre prend effet uniquement lorsque le paramètre sink.sync-write est défini sur true.
sink.ignore-null
Indique s'il faut ignorer les valeurs null.
Boolean
No
false
RemarqueSi ce paramètre est défini sur true, le paramètre
null-string-literalne prend pas effet.Ce paramètre est pris en charge uniquement par Realtime Compute for Apache Flink version VVR 8.0.9 ou ultérieure.
-
Options spécifiques aux dimensions (liées au cache)
Option
Description
Type de données
Obligatoire
Valeur par défaut
Remarques
cache
La politique de mise en cache.
String
Non
ALL
Valeurs valides :
None : aucune donnée n'est mise en cache.
LRU : seules certaines données de la table de dimension sont mises en cache. À chaque réception d'un enregistrement, le système interroge le cache. Si l'enregistrement est absent du cache, le système recherche l'enregistrement dans la table de dimension physique.
RemarqueSi vous utilisez cette politique de mise en cache, vous devez configurer les paramètres cacheSize et cacheTTLMs.
ALL : toutes les données de la table de dimension sont mises en cache. Il s'agit de la valeur par défaut. Avant l'exécution d'une tâche, le système charge l'intégralité des données de la table de dimension dans le cache. Ainsi, toutes les requêtes ultérieures dans la table de dimension sont résolues via le cache. Si le système ne trouve pas l'enregistrement dans le cache, cela signifie que la clé de jointure n'existe pas. Le système recharge l'ensemble des données du cache après l'expiration des entrées.
RemarqueSi le volume de données d'une table distante est faible et qu'un grand nombre de clés manquantes existent, nous vous recommandons de définir ce paramètre sur ALL. La table source et la table de dimension ne peuvent pas être associées selon la clause ON. Si vous utilisez cette politique de mise en cache, vous devez configurer les paramètres cacheTTLMs et cacheReloadTimeBlackList.
Si toutes les données de la table de dimension sont chargées dans le cache, la vitesse de démarrage du déploiement peut ralentir. Vous pouvez configurer la politique de mise en cache de manière flexible en fonction de vos besoins métier.
Si vous définissez le paramètre cache sur ALL, vous devez augmenter la mémoire du nœud utilisé pour la jointure des tables, car le système charge de manière asynchrone les données depuis la table de dimension. La taille mémoire supplémentaire correspond au double de celle de la table distante.
cacheSize
Le nombre maximal de lignes de données pouvant être mises en cache.
Long
Non
10000
Vous pouvez configurer ce paramètre lorsque vous définissez le paramètre cache sur LRU.
cacheTTLMs
La durée d'expiration du cache. Unité : millisecondes.
Long
Non
Aucune valeur par défaut
La configuration du paramètre cacheTTLMs varie en fonction du paramètre cache.
Si vous définissez le paramètre cache sur None, vous pouvez laisser le paramètre cacheTTLMs vide. Cela indique que les entrées du cache n'expirent pas.
Si vous définissez le paramètre cache sur LRU, le paramètre cacheTTLMs spécifie la durée d'expiration du cache. Par défaut, les entrées du cache n'expirent pas.
Si vous définissez le paramètre cache sur ALL, le paramètre cacheTTLMs spécifie l'intervalle auquel le système recharge le cache. Par défaut, le cache n'est pas rechargé.
cacheEmpty
Indique s'il faut mettre en cache les résultats vides.
Boolean
Non
true
S.O.
cacheReloadTimeBlackList
Les plages horaires durant lesquelles le cache n'est pas actualisé. Ce paramètre prend effet lorsque le paramètre cache est défini sur ALL. Le cache n'est pas actualisé pendant les périodes que vous spécifiez pour ce paramètre. Ce paramètre est adapté aux événements promotionnels en ligne à grande échelle, tels que le Double 11.
String
Non
Aucune valeur par défaut
L'exemple suivant illustre le format des valeurs : 2017-10-24 14:00 -> 2017-10-24 15:00, 2017-11-10 23:30 -> 2017-11-11 08:00. Utilisez des délimiteurs selon les règles suivantes :
Séparez plusieurs plages horaires par des virgules (,).
Séparez l'heure de début et l'heure de fin de chaque plage horaire par une flèche (->) constituée d'un trait d'union (-) et d'un chevron fermant (>).
cacheScanLimit
Le nombre de lignes renvoyées par le serveur RPC (Remote Procedure Call) à un client lorsque le serveur lit l'intégralité des données d'une table de dimension ApsaraDB for HBase.
Integer
Non
100
Ce paramètre est disponible uniquement lorsque vous définissez le paramètre cache sur ALL.
Mappages des types de données
Dans une table ApsaraDB for HBase, les valeurs des types de données de Realtime Compute for Apache Flink sont converties en tableaux d'octets à l'aide de org.apache.hadoop.hbase.util.Bytes. Le processus de décodage varie selon les scénarios suivants :
Si le type de données de Realtime Compute for Apache Flink n'est pas STRING et qu'une valeur dans la table ApsaraDB for HBase est un tableau d'octets vide, la valeur est décodée comme null.
Si le type de données de Realtime Compute for Apache Flink est STRING et qu'une valeur dans la table de dimension ApsaraDB for HBase correspond au tableau d'octets spécifié par
null-string-literal, la valeur est décodée comme null.
Type SQL Flink | Fonction utilisée pour convertir une valeur en octets pour ApsaraDB for HBase | Fonction utilisée pour lire les octets depuis ApsaraDB for HBase |
CHAR | byte[] toBytes(String s) | String toString(byte[] b) |
VARCHAR | ||
STRING | ||
BOOLEAN | byte[] toBytes(boolean b) | boolean toBoolean(byte[] b) |
BINARY | byte[] | byte[] |
VARBINARY | ||
DECIMAL | byte[] toBytes(BigDecimal v) | BigDecimal toBigDecimal(byte[] b) |
TINYINT | new byte[] { val } | bytes[0] |
SMALLINT | byte[] toBytes(short val) | short toShort(byte[] bytes) |
INT | byte[] toBytes(int val) | int toInt(byte[] bytes) |
BIGINT | byte[] toBytes(long val) | long toLong(byte[] bytes) |
FLOAT | byte[] toBytes(float val) | float toFloat(byte[] bytes) |
DOUBLE | byte[] toBytes(double val) | double toDouble(byte[] bytes) |
DATE | Convertit une date en valeur INT représentant le nombre de jours écoulés depuis le 1er janvier 1970, puis en tableau d'octets à l'aide de | Convertit un tableau d'octets de la base de données ApsaraDB for HBase en type de données INT à l'aide de |
TIME | Convertit une heure en valeur INT représentant le nombre de millisecondes écoulées depuis 00:00:00, puis en tableau d'octets à l'aide de | Convertit un tableau d'octets de la base de données ApsaraDB for HBase en type de données INT à l'aide de |
TIMESTAMP | Convertit un horodatage en valeur LONG représentant le nombre de millisecondes écoulées depuis le 1er janvier 1970 à 00:00:00, puis en tableau d'octets à l'aide de | Convertit un tableau d'octets de la base de données ApsaraDB for HBase en type de données LONG à l'aide de |
Exemples de code
-
Exemple de code pour une table de dimension
CREATE TEMPORARY TABLE datagen_source ( a INT, b BIGINT, c STRING, `proc_time` AS PROCTIME() ) WITH ( 'connector'='datagen' ); CREATE TEMPORARY TABLE hbase_dim ( rowkey INT, family1 ROW<col1 INT>, family2 ROW<col1 STRING, col2 BIGINT>, family3 ROW<col1 DOUBLE, col2 BOOLEAN, col3 STRING> ) WITH ( 'connector' = 'cloudhbase', 'table-name' = '<yourTableName>', 'zookeeper.quorum' = '<yourZookeeperQuorum>' ); CREATE TEMPORARY TABLE blackhole_sink( a INT, f1c1 INT, f3c3 STRING ) WITH ( 'connector' = 'blackhole' ); INSERT INTO blackhole_sink SELECT a, family1.col1 as f1c1, family3.col3 as f3c3 FROM datagen_source JOIN hbase_dim FOR SYSTEM_TIME AS OF datagen_source.`proc_time` as h ON datagen_source.a = h.rowkey; -
Exemple de code pour une table de destination
CREATE TEMPORARY TABLE datagen_source ( rowkey INT, f1q1 INT, f2q1 STRING, f2q2 BIGINT, f3q1 DOUBLE, f3q2 BOOLEAN, f3q3 STRING ) WITH ( 'connector'='datagen' ); CREATE TEMPORARY TABLE hbase_sink ( rowkey INT, family1 ROW<q1 INT>, family2 ROW<q1 STRING, q2 BIGINT>, family3 ROW<q1 DOUBLE, q2 BOOLEAN, q3 STRING>, PRIMARY KEY (rowkey) NOT ENFORCED ) WITH ( 'connector'='cloudhbase', 'table-name'='<yourTableName>', 'zookeeper.quorum'='<yourZookeeperQuorum>' ); INSERT INTO hbase_sink SELECT rowkey, ROW(f1q1), ROW(f2q1, f2q2), ROW(f3q1, f3q2, f3q3) FROM datagen_source; -
Exemple de code pour une table de destination prenant en charge les colonnes dynamiques
CREATE TEMPORARY TABLE datagen_source ( id INT, f1hour STRING, f1deal BIGINT, f2day STRING, f2deal BIGINT ) WITH ( 'connector'='datagen' ); CREATE TEMPORARY TABLE hbase_sink ( rowkey INT, f1 ROW<`hour` STRING, deal BIGINT>, f2 ROW<`day` STRING, deal BIGINT> ) WITH ( 'connector'='cloudhbase', 'table-name'='<yourTableName>', 'zookeeper.quorum'='<yourZookeeperQuorum>', 'dynamic.table'='true' ); INSERT INTO hbase_sink SELECT id, ROW(f1hour, f1deal), ROW(f2day, f2deal) FROM datagen_source;Si dynamic.table est défini sur true, une table ApsaraDB for HBase prenant en charge les colonnes dynamiques est utilisée.
Deux champs doivent être déclarés dans les lignes correspondant à chaque famille de colonnes. La valeur du premier champ indique la colonne dynamique, et la valeur du second champ indique la valeur de la colonne dynamique.
Par exemple, la table datagen_source contient une ligne de données. Cette ligne indique que l'ID du produit est 1, le montant de la transaction du produit entre 10:00 et 11:00 est de 100, et le montant de la transaction du produit le 26 juillet 2020 est de 10000. Dans ce cas, une ligne dont la rowkey est 1 est insérée dans la table ApsaraDB for HBase. f1:10 vaut 100, et f2:2020-7-26 vaut 10000.