Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector do ApsaraDB for HBase

Última atualização: Jun 27, 2026

Este tópico descreve como usar o conector do ApsaraDB for HBase.

Informações básicas

O ApsaraDB for HBase é um serviço NoSQL inteligente baseado em nuvem e com excelente custo-benefício. Ele oferece alta escalabilidade e compatibilidade com o HBase open source. O ApsaraDB for HBase proporciona vantagens como baixos custos de armazenamento, alto throughput, escalabilidade e processamento inteligente de dados. Esse serviço suporta funcionalidades essenciais da Alibaba, incluindo recomendações do Taobao, controle de risco do Ant Credit Pay, publicidade, painéis de dados, rastreamento logístico da Cainiao, registros de transações do Alipay e mensagens do Taobao Mobile. Trata-se de um serviço totalmente gerenciado, com recursos de nível empresarial como processamento de petabytes de dados, alta concorrência, dimensionamento rápido em segundos, latência de resposta na casa dos milissegundos, alta disponibilidade entre data centers e distribuição global.

A tabela a seguir descreve as capacidades suportadas pelo conector do ApsaraDB for HBase.

Item

Descrição

Tipo de tabela

Tabela de dimensão e tabela sink

Modo de execução

Modo streaming

Formato de dados

N/A

Métrica

Métricas

  • source

    Nenhuma

  • Tabelas de dimensão

    Nenhuma

  • Sink

    numBytesOut, numBytesOutPerSecond, numRecordsOut, numRecordsOutPerSecond e currentSendTime

    Nota

    Para mais informações sobre as métricas, consulte Métricas.

Tipo de API

API SQL

Atualização ou exclusão de dados em uma tabela sink

Suportado

Pré-requisitos

  • Adquira um cluster do ApsaraDB for HBase e crie uma tabela do ApsaraDB for HBase. Para mais detalhes sobre como adquirir um cluster do ApsaraDB for HBase, consulte Adquirir um cluster.

  • Configure uma lista de permissões para o cluster do ApsaraDB for HBase. Para mais informações, consulte Configurar uma lista de permissões.

Observações de uso

Antes de usar o conector do ApsaraDB for HBase, confirme o tipo da sua instância de banco de dados e certifique-se de que o tipo de conector selecionado está correto. O uso inadequado do conector pode causar problemas inesperados.

  • O conector descrito neste tópico destina-se a instâncias do ApsaraDB for HBase.

  • Instâncias do Lindorm são compatíveis com Apache HBase. Para instâncias Lindorm, use o conector específico do Lindorm. Para mais informações, consulte Lindorm.

  • Se você usar o conector do ApsaraDB for HBase para conectar o Realtime Compute for Apache Flink a um banco de dados HBase open source, a validade dos dados não poderá ser garantida.

Sintaxe

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>'
);
  • Declare as famílias de colunas de uma tabela do ApsaraDB for HBase com o tipo ROW. Cada nome de família de colunas corresponde ao nome de um campo da linha. Na sintaxe DDL, as seguintes famílias de colunas são declaradas: family1, family2 e family3.

  • Uma coluna dentro de uma família de colunas corresponde a um campo em uma linha. O nome da coluna é igual ao nome do campo. Na sintaxe DDL, as colunas q2 e q3 são declaradas na família de colunas family2.

  • Além dos campos do tipo ROW, apenas um campo de tipo atômico, como STRING ou BIGINT, pode existir em uma tabela do ApsaraDB for HBase. Esse campo atômico é considerado a chave de linha (row key) da tabela, como o campo rowkey na instrução DDL.

  • Defina a chave de linha de uma tabela do ApsaraDB for HBase como a chave primária da tabela sink. Se nenhuma chave primária for definida, a chave de linha será usada automaticamente como chave primária.

  • Declare na tabela sink apenas as famílias de colunas e colunas necessárias da tabela do ApsaraDB for HBase.

Opções do conector

  • Gerais

    Opção

    Descrição

    Tipo de dado

    Obrigatório

    Valor padrão

    Observações

    connector

    O tipo da tabela.

    String

    Sim

    Sem valor padrão

    Defina o valor como cloudhbase.

    table-name

    O nome da tabela do ApsaraDB for HBase.

    String

    Sim

    Sem valor padrão

    N/A.

    zookeeper.znode.quorum

    A URL usada para acessar o serviço ZooKeeper do ApsaraDB for HBase.

    String

    Sim

    Sem valor padrão

    N/A.

    zookeeper.znode.parent

    O diretório raiz do ApsaraDB for HBase no serviço ZooKeeper.

    String

    Não

    /hbase

    Este parâmetro tem efeito apenas no ApsaraDB for HBase Standard Edition.

    userName

    O nome de usuário usado para acessar o banco de dados.

    String

    Não

    Sem valor padrão

    Este parâmetro tem efeito apenas no ApsaraDB for HBase Performance-enhanced Edition.

    password

    A senha usada para acessar o banco de dados.

    String

    Não

    Sem valor padrão

    Este parâmetro tem efeito apenas no ApsaraDB for HBase Performance-enhanced Edition.

    haclient.cluster.id

    O ID do cluster do ApsaraDB for HBase em modo de alta disponibilidade (HA).

    String

    Não

    Sem valor padrão

    Este parâmetro é necessário apenas ao acessar clusters com recuperação de desastres entre zonas. Tem efeito apenas no ApsaraDB for HBase Performance-enhanced Edition.

    retires.number

    O número de tentativas permitidas para o cliente do ApsaraDB for HBase se conectar ao banco de dados.

    Integer

    Não

    31

    N/A.

    null-string-literal

    Se o tipo de dado de um campo no ApsaraDB for HBase for STRING e o dado correspondente no Realtime Compute for Apache Flink for nulo, o valor de null-string-literal será atribuído a esse campo e gravado no banco de dados.

    String

    Não

    null

    N/A.

  • Específicas para sink

    Opção

    Descrição

    Tipo de dado

    Obrigatório

    Valor padrão

    Observações

    sink.buffer-flush.max-size

    O tamanho dos dados em bytes armazenados em cache na memória antes da gravação no banco de dados ApsaraDB for HBase. Um valor maior melhora a performance de gravação, mas aumenta a latência e o consumo de memória.

    String

    Não

    2MB

    Unidade: B, KB, MB ou GB (não diferencia maiúsculas de minúsculas). Se definido como 0, nenhum dado será armazenado em cache.

    sink.buffer-flush.max-rows

    O número de registros de dados armazenados em cache na memória antes da gravação no banco de dados ApsaraDB for HBase. Valores maiores melhoram a performance de gravação, porém aumentam a latência e o consumo de memória.

    Integer

    Não

    1000

    Se definido como 0, nenhum dado será armazenado em cache.

    sink.buffer-flush.interval

    O intervalo no qual os dados em cache são gravados no banco de dados ApsaraDB for HBase. Este parâmetro controla a latência de gravação dos dados.

    Duration

    Não

    1s

    Unidade: ms, s, min, h ou d. Se definido como 0, a gravação periódica de dados é desativada.

    dynamic.table

    Define se deve ser usada uma tabela do ApsaraDB for HBase com suporte a colunas dinâmicas.

    Boolean

    Não

    false

    Valores válidos:

    • true

    • false

    sink.ignore-delete

    Define se as mensagens de retração devem ser ignoradas.

    Boolean

    Não

    false

    Se um stream contiver eventos DELETE ou UPDATE_BEFORE e múltiplas tarefas sink atualizarem concorrentemente diferentes campos de uma tabela, pode ocorrer inconsistência de dados.

    Por exemplo, após a exclusão de um registro, outra tarefa pode atualizar alguns campos. Os campos não atualizados tornar-se-ão nulos ou receberão valores padrão, causando erros nos dados.

    Para evitar esse problema, defina sink.ignore-delete como true para ignorar eventos upstream de DELETE e UPDATE_BEFORE.

    Nota
    • UPDATE_BEFORE faz parte do mecanismo de retração do Flink e serve para retrair o valor antigo durante uma operação de atualização.

    • Se ignoreDelete estiver definido como true, todos os eventos DELETE e UPDATE_BEFORE serão ignorados. Apenas registros INSERT e UPDATE_AFTER serão processados.

    sink.sync-write

    Define se os dados devem ser gravados no ApsaraDB for HBase em modo síncrono.

    Boolean

    Não

    true

    Valores válidos:

    • true: Os dados são gravados em modo síncrono. Nesse modo, a gravação ocorre em sequência, mas a performance é reduzida.

    • false: Os dados são gravados em modo assíncrono. Nesse modo, a ordem de gravação pode não ser preservada, mas a performance é melhorada.

    sink.buffer-flush.batch-rows

    O número de registros de dados armazenados em cache na memória quando a gravação no ApsaraDB for HBase ocorre em modo síncrono. Valores maiores melhoram a performance de gravação, mas aumentam a latência e o uso de memória.

    Integer

    Não

    100

    Este parâmetro tem efeito apenas quando o parâmetro sink.sync-write está definido como true.

    sink.ignore-null

    Define se valores nulos devem ser ignorados.

    Boolean

    Não

    false

    Nota
    • Se este parâmetro estiver definido como true, o parâmetro null-string-literal não terá efeito.

    • Apenas o Realtime Compute for Apache Flink com VVR 8.0.9 ou superior suporta este parâmetro.

  • Opções específicas para tabela de dimensão (relacionadas a cache)

    Opção

    Descrição

    Tipo de dado

    Obrigatório

    Valor padrão

    Observações

    cache

    A política de cache.

    String

    Não

    ALL

    Valores válidos:

    • None: Nenhum dado é armazenado em cache.

    • LRU: Apenas dados específicos da tabela de dimensão são armazenados em cache. Sempre que o sistema recebe um registro de dados, ele busca no cache. Caso não encontre o registro, o sistema consulta a tabela de dimensão física.

      Nota

      Ao usar esta política de cache, configure os parâmetros cacheSize e cacheTTLMs.

    • ALL: Todos os dados da tabela de dimensão são armazenados em cache. Este é o valor padrão. Antes da execução de um job, o sistema carrega todos os dados da tabela de dimensão para o cache. Assim, todas as consultas subsequentes buscam primeiro no cache. Se o registro não for encontrado no cache, considera-se que a chave de junção não existe. O sistema recarrega todos os dados no cache após a expiração das entradas.

      Nota
      • Se o volume de dados na tabela remota for pequeno e houver muitas chaves ausentes, recomenda-se definir este parâmetro como ALL. A tabela de origem e a tabela de dimensão não podem ser associadas pela cláusula ON. Ao usar esta política, configure os parâmetros cacheTTLMs e cacheReloadTimeBlackList.

      • Carregar todos os dados da tabela de dimensão para o cache pode reduzir a velocidade de inicialização do deployment. Configure a política de cache flexivelmente conforme suas necessidades de negócio.

    Ao definir o parâmetro cache como ALL, aumente a memória do nó responsável pela junção das tabelas, pois o sistema carrega dados da tabela de dimensão de forma assíncrona. O aumento de memória necessário equivale ao dobro do tamanho da tabela remota.

    cacheSize

    O número máximo de linhas de dados que podem ser armazenadas em cache.

    Long

    Não

    10000

    Configure este parâmetro quando definir o parâmetro cache como LRU.

    cacheTTLMs

    O tempo limite do cache. Unidade: milissegundos.

    Long

    Não

    Sem valor padrão

    A configuração do parâmetro cacheTTLMs varia conforme o parâmetro cache.

    • Se o parâmetro cache for definido como None, o parâmetro cacheTTLMs pode ficar vazio, indicando que as entradas de cache não expiram.

    • Se o parâmetro cache for definido como LRU, o parâmetro cacheTTLMs especifica o tempo limite do cache. Por padrão, as entradas não expiram.

    • Se o parâmetro cache for definido como ALL, o parâmetro cacheTTLMs especifica o intervalo no qual o sistema recarrega o cache. Por padrão, o cache não é recarregado.

    cacheEmpty

    Define se resultados vazios devem ser armazenados em cache.

    Boolean

    Não

    true

    N/A.

    cacheReloadTimeBlackList

    Os períodos durante os quais o cache não é atualizado. Este parâmetro tem efeito quando o parâmetro cache está definido como ALL. O cache não será atualizado nos intervalos especificados, sendo ideal para grandes eventos promocionais online, como o Double 11.

    String

    Não

    Sem valor padrão

    Exemplo de formato de valores: 2017-10-24 14:00 -> 2017-10-24 15:00, 2017-11-10 23:30 -> 2017-11-11 08:00. Use delimitadores conforme as regras abaixo:

    • Separe múltiplos períodos de tempo com vírgulas (,).

    • Separe a hora de início e a hora de fim de cada período com uma seta (->), formada por um hífen (-) e um sinal de maior que (>).

    cacheScanLimit

    O número de linhas que o servidor de chamada de procedimento remoto (RPC) retorna ao cliente ao ler todos os dados de uma tabela de dimensão do ApsaraDB for HBase.

    Integer

    Não

    100

    Este parâmetro está disponível apenas quando o parâmetro cache está definido como ALL.

Mapeamentos de tipos de dados

Um valor de um tipo de dado do Realtime Compute for Apache Flink é convertido em um array de bytes usando org.apache.hadoop.hbase.util.Bytes em uma tabela do ApsaraDB for HBase. O processo de decodificação varia conforme os cenários abaixo:

  • Se o tipo de dado do Realtime Compute for Apache Flink não for STRING e o valor na tabela do ApsaraDB for HBase for um array de bytes vazio, o valor será decodificado como nulo.

  • Se o tipo de dado do Realtime Compute for Apache Flink for STRING e o valor na tabela de dimensão do ApsaraDB for HBase for o array de bytes especificado por null-string-literal, o valor será decodificado como nulo.

Tipo Flink SQL

Função para converter valor em bytes para o ApsaraDB for HBase

Função para ler bytes do 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

Converte uma data em um valor INT representando o número de dias desde 1º de janeiro de 1970 e, em seguida, em um array de bytes usando byte[] toBytes(int val).

Converte um array de bytes do banco de dados ApsaraDB for HBase no tipo de dado INT usando int toInt(byte[] bytes). O valor INT representa o número de dias desde 1º de janeiro de 1970.

TIME

Converte um horário em um valor INT representando o número de milissegundos desde 00:00:00 e, em seguida, em um array de bytes usando byte[] toBytes(int val).

Converte um array de bytes do banco de dados ApsaraDB for HBase no tipo de dado INT usando int toInt(byte[] bytes). O valor INT representa o número de milissegundos desde 00:00:00.

TIMESTAMP

Converte um timestamp em um valor LONG representando o número de milissegundos desde 00:00:00 de 1º de janeiro de 1970 e, em seguida, em um array de bytes usando byte[] toBytes(long val).

Converte um array de bytes do banco de dados ApsaraDB for HBase no tipo de dado LONG usando long toLong(byte[] bytes). O valor LONG representa o número de milissegundos desde 00:00:00 de 1º de janeiro de 1970.

Código de exemplo

  • Código de exemplo para tabela de dimensão

    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;
  • Código de exemplo para tabela sink

    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;
  • Código de exemplo para tabela sink com suporte a colunas dinâmicas

    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;
    • Se dynamic.table estiver definido como true, será usada uma tabela do ApsaraDB for HBase com suporte a colunas dinâmicas.

    • Declare dois campos nas linhas correspondentes a cada família de colunas. O valor do primeiro campo indica a coluna dinâmica, enquanto o valor do segundo campo indica o valor dessa coluna dinâmica.

    • Por exemplo, a tabela datagen_source contém uma linha de dados indicando que o ID da mercadoria é 1, o valor da transação entre 10:00 e 11:00 é 100, e o valor da transação em 26 de julho de 2020 é 10000. Nesse caso, uma linha com rowkey igual a 1 é inserida na tabela do ApsaraDB for HBase. f1:10 será 100 e f2:2020-7-26 será 10000.