Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:ApsaraDB for HBase connector

Dernière mise à jour :Aug 09, 2026

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

Métriques

  • Source

    Aucune

  • Tables de dimension

    Aucune

  • Destination (Sink)

    numBytesOut, numBytesOutPerSecond, numRecordsOut, numRecordsOutPerSecond et currentSendTime

    Remarque

    Pour plus d'informations sur les métriques, consultez Métriques.

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 rowkey dans 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

    /hbase

    Ce 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-literal est 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 DELETE ou UPDATE_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-delete sur true afin d'ignorer les événements DELETE et UPDATE_BEFORE en amont.

    Remarque
    • UPDATE_BEFORE fait 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 ignoreDelete est défini sur true, tous les événements DELETE et UPDATE_BEFORE sont ignorés. Seuls les enregistrements INSERT et UPDATE_AFTER sont 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

    Remarque
    • Si ce paramètre est défini sur true, le paramètre null-string-literal ne 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.

      Remarque

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

      Remarque
      • Si 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 byte[] toBytes(int val).

Convertit un tableau d'octets de la base de données ApsaraDB for HBase en type de données INT à l'aide de int toInt(byte[] bytes). La valeur du type de données INT représente le nombre de jours écoulés depuis le 1er janvier 1970.

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 byte[] toBytes(int val).

Convertit un tableau d'octets de la base de données ApsaraDB for HBase en type de données INT à l'aide de int toInt(byte[] bytes). La valeur du type de données INT représente le nombre de millisecondes écoulées depuis 00:00:00.

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 byte[] toBytes(long val).

Convertit un tableau d'octets de la base de données ApsaraDB for HBase en type de données LONG à l'aide de long toLong(byte[] bytes). La valeur du type de données LONG représente le nombre de millisecondes écoulées depuis le 1er janvier 1970 à 00:00:00.

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.