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 |
|
|
Nome da tabela de dimensão do Lindorm |
|
|
Nome da tabela sink do Lindorm |
|
|
Endpoint do servidor Lindorm no formato |
|
|
Nome de usuário do Lindorm |
|
|
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 |
|
|
String |
Sim |
— |
Use |
|
|
String |
Sim |
— |
Endpoint do servidor Lindorm no formato |
|
|
String |
Sim |
— |
Namespace do banco de dados Lindorm. |
|
|
String |
Sim |
— |
Nome de usuário do banco de dados Lindorm. |
|
|
String |
Sim |
— |
Senha do banco de dados Lindorm. |
|
|
String |
Sim |
— |
Nome da tabela do Lindorm. |
|
|
String |
Sim |
— |
Nome da família de colunas. Se nenhuma família de colunas foi especificada durante a criação da tabela, insira |
|
|
Integer |
Não |
|
Intervalo entre novas tentativas para operações de leitura com falha, em milissegundos. |
|
|
Integer |
Não |
|
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 |
|
|
Integer |
Não |
|
Quantidade de registros armazenados em buffer antes do flush para o Lindorm. |
|
|
Integer |
Não |
|
Tempo máximo entre flushes quando o buffer não está cheio, em milissegundos. |
|
|
Boolean |
Não |
|
Se |
|
|
Boolean |
Não |
|
Se |
|
|
String |
Não |
— |
Lista separada por vírgulas das colunas a excluir das atualizações. Por exemplo, |
Específico para tabela de dimensão
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
Boolean |
Não |
|
Se |
|
|
Boolean |
Não |
|
Se |
|
|
String |
Não |
|
Política de cache. Valores válidos: |
|
|
Integer |
Não |
|
Número máximo de linhas mantidas em cache. Aplica-se quando |
|
|
Integer |
Não |
— |
Tempo de expiração das entradas de cache, em milissegundos. Aplica-se quando |
|
|
Boolean |
Não |
|
Se |
|
|
Boolean |
Não |
|
Se |
|
|
Integer |
Não |
|
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
cacheSizelinhas (as linhas mais antigas são removidas primeiro).A entrada permanece no cache por mais tempo que o definido em
cacheTTLMsmilissegundos.
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
c1ec2.
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 |
|
|
|
|
|
|
|
|
|
|
Bytes diretos |
Bytes diretos |
|
|
|
|
|
|
Primeiro byte de |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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 |
|
|
Total de bytes gravados no sink |
|
|
Bytes gravados por segundo |
|
|
Total de registros gravados no sink |
|
|
Registros gravados por segundo |