Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Lindorm connector

Última atualização: Jun 27, 2026

Sink: Streaming Lookup source: Sync/Async mode

O conector do Lindorm permite que jobs de streaming do Flink gravem e consultem dados em tabelas amplas do Lindorm pela API SQL.

Pré-requisitos

Antes de começar, verifique se você tem:

  • Um mecanismo de tabela ampla do Lindorm e uma tabela do Lindorm. Para mais informações, consulte Crie uma instância.

  • Conectividade de rede entre o cluster do Lindorm e o workspace do Flink (por exemplo, ambos na mesma Virtual Private Cloud (VPC)).

Não há suporte para tabelas HBase do Lindorm. Apenas o LindormTable é compatível.

Início rápido

O exemplo a seguir gera 10 linhas de dados, consulta linhas correspondentes em uma tabela de dimensão do Lindorm e grava o resultado em uma tabela sink do 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;

Substitua os placeholders antes de executar:

Placeholder

Descrição

${lindorm_dim_table}

Nome da tabela de dimensão do Lindorm

${lindorm_sink_table}

Nome da tabela sink do Lindorm

${lindorm_seed_server}

Endpoint do servidor Lindorm no formato host:port

${lindorm_username}

Nome de usuário do Lindorm

${lindorm_password}

Senha do Lindorm

Sintaxe

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>'
);

Opções do conector

Geral

Opção

Tipo

Obrigatório

Padrão

Descrição

connector

String

Sim

Use lindorm.

seedserver

String

Sim

Endpoint do servidor Lindorm no formato host:port. O Realtime Compute for Apache Flink usa a API ApsaraDB for HBase para Java para se conectar. Para obter detalhes, consulte Use Flink to connect to and use LindormTable.

namespace

String

Sim

Namespace do banco de dados Lindorm.

username

String

Sim

Nome de usuário do banco de dados Lindorm.

password

String

Sim

Senha do banco de dados Lindorm.

tableName

String

Sim

Nome da tabela do Lindorm.

columnFamily

String

Sim

Nome da família de colunas. Se nenhuma família de colunas foi especificada durante a criação da tabela, insira f.

retryIntervalMs

Integer

Não

1000

Intervalo entre novas tentativas para operações de leitura com falha, em milissegundos.

maxRetryTimes

Integer

Não

5

Número máximo de novas tentativas para operações de leitura ou gravação.

Específico para sink

Opção

Tipo

Obrigatório

Padrão

Descrição

bufferSize

Integer

Não

500

Quantidade de registros armazenados em buffer antes do flush para o Lindorm.

flushIntervalMs

Integer

Não

2000

Tempo máximo entre flushes quando o buffer não está cheio, em milissegundos.

ignoreDelete

Boolean

Não

false

Se true, ignora as operações de exclusão.

dynamicColumnSink

Boolean

Não

false

Se true, habilita o recurso de tabela dinâmica. Consulte Tabela dinâmica.

excludeUpdateColumns

String

Não

Lista separada por vírgulas das colunas a excluir das atualizações. Por exemplo, a,b,c ignora atualizações nas colunas a, b e c. Requer VVR 8.0.9 ou posterior.

Específico para tabela de dimensão

Opção

Tipo

Obrigatório

Padrão

Descrição

partitionedJoin

Boolean

Não

false

Se true, usa a JoinKey para particionamento e aumenta a taxa de acerto do cache.

shuffleEmptyKey

Boolean

Não

false

Se true, encaminha chaves vazias do upstream aleatoriamente para nós downstream. Se false, direciona-as para a thread paralela 0.

cache

String

Não

None

Política de cache. Valores válidos: None (sem cache) e LRU (armazena em cache as linhas acessadas recentemente).

cacheSize

Integer

Não

1000

Número máximo de linhas mantidas em cache. Aplica-se quando cache é LRU.

cacheTTLMs

Integer

Não

Tempo de expiração das entradas de cache, em milissegundos. Aplica-se quando cache é LRU. Por padrão, as entradas não expiram.

cacheEmpty

Boolean

Não

true

Se true, armazena em cache resultados de consulta sem retorno de linhas.

async

Boolean

Não

false

Se true, habilita o modo de consulta assíncrona.

asyncLindormRpcTimeoutMs

Integer

Não

300000

Timeout de RPC para consultas assíncronas, em milissegundos.

Cache de consulta

Por padrão, o conector busca cada consulta diretamente no Lindorm (cache = None). Ative o cache LRU para reduzir a pressão de leitura no Lindorm em joins de alto throughput.

Com cache = LRU, o conector armazena as linhas acessadas recentemente na memória. Quando a chave de join corresponde a uma linha em cache, o conector retorna o valor armazenado sem consultar o Lindorm. Uma entrada em cache é removida quando:

  • O cache atinge cacheSize linhas (as linhas mais antigas são removidas primeiro).

  • A entrada permanece no cache por mais tempo que o definido em cacheTTLMs milissegundos.

Há um compromisso: um TTL maior ou um cache mais amplo reduz o tráfego de leitura do Lindorm, mas aumenta o risco de servir dados desatualizados. Ajuste ambos os valores com base nos requisitos de throughput e na frequência de alteração dos dados subjacentes.

O conector do Lindorm suporta joins de consulta um-para-muitos. Preste atenção às estratégias de cache e ao throughput quando as linhas da tabela de dimensão puderem corresponder a múltiplos eventos do upstream.

Gravações idempotentes

Quando uma tabela do Lindorm possui chave primária, todas as gravações usam semântica upsert: cada registro recebido insere uma nova linha ou atualiza a linha existente correspondente. Isso torna as gravações idempotentes.

Gravações idempotentes são fundamentais para a tolerância a falhas. Se um job do Flink reiniciar a partir de um checkpoint, ele reprocessará as mensagens desde o último checkpoint bem-sucedido. Como os upserts do Lindorm são idempotentes, os registros reprocessados produzem o mesmo resultado das gravações originais, sem linhas duplicadas ou violações de restrições.

Defina uma chave primária na DDL para aproveitar as gravações idempotentes.

Tabela dinâmica

Use o recurso de tabela dinâmica quando o schema evoluir em tempo de execução. Nesse cenário, as colunas são criadas dinamicamente com base nos valores dos dados, em vez de serem fixadas na DDL. Um caso de uso típico é o rastreamento de métricas horárias por dia, em que as horas são nomes de colunas e os dias são chaves primárias:

Chave primária

00:00

01:00

2025-06-01

45

32

2025-06-02

76

34

Regras de DDL para tabelas dinâmicas:

  • As primeiras N colunas formam a chave primária.

  • As duas últimas colunas devem ser do tipo VARCHAR.

  • A penúltima coluna (c1) contém o nome da coluna dinâmica.

  • A última coluna (c2) armazena o valor dessa coluna.

  • Não são permitidas colunas fora da chave primária além de c1 e c2.

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'
);

Sempre que um registro é gravado, o conector adiciona ou atualiza uma coluna na linha do Lindorm identificada por <pk1, pk2, pk3>. As demais colunas dessa linha permanecem inalteradas.

Mapeamentos de tipos de dados

Todos os dados do Lindorm são armazenados em formato binário. A tabela a seguir mostra como o conector converte os tipos do Flink SQL para as representações binárias do Lindorm e vice-versa.

Tipo Flink SQL

Gravar no Lindorm

Ler do Lindorm

CHAR / VARCHAR

StringData::toBytes

StringData::fromBytes

BOOLEAN

Bytes::toBytes(boolean)

Bytes::toBigDecimal

BINARY / VARBINARY

Bytes diretos

Bytes diretos

DECIMAL

Bytes::toBytes(BigDecimal)

Bytes::toBigDecimal

TINYINT

Primeiro byte 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) com dias desde 1970-01-01

Bytes::toInt → dias desde 1970-01-01

TIME

Bytes::toBytes(int) com milissegundos desde 00:00:00

Bytes::toInt → milissegundos desde 00:00:00

TIMESTAMP

Bytes::toBytes(long) com milissegundos desde 1970-01-01 00:00:00

Bytes::toLong → milissegundos desde 1970-01-01 00:00:00

Os métodos de Bytes estão na classe com.alibaba.lindorm.client.core.utils.Bytes. Os métodos de StringData (linhas CHAR/VARCHAR) encontram-se na classe org.apache.flink.table.data.StringData.

Métricas

As seguintes métricas de sink estão disponíveis. Para obter detalhes, consulte Metrics.

Métrica

Descrição

numBytesOut

Total de bytes gravados no sink

numBytesOutPerSecond

Bytes gravados por segundo

numRecordsOut

Total de registros gravados no sink

numRecordsOutPerSecond

Registros gravados por segundo

FAQ

Erros de conexão do Lindorm e soluções