O conector do Tair (Enterprise Edition) grava dados de streaming do Flink em uma instância do Tair (Enterprise Edition). O Tair é um banco de dados compatível com Redis que oferece armazenamento híbrido em memória e disco, standby ativo para alta disponibilidade e arquitetura de cluster escalável para cargas de trabalho com alto throughput e baixa latência.
Recursos suportados
|
Categoria |
Detalhes |
|
Tipos suportados |
Tabela sink |
|
Modo de execução |
Modo stream |
|
Formato de dados |
STRING |
|
Tipo de API |
SQL |
|
Atualização ou exclusão de dados na tabela sink |
Sim |
|
Métricas exclusivas |
|
Para obter detalhes sobre as métricas, consulte Métricas.
Pré-requisitos
Antes de começar, verifique se você possui:
Uma instância do Tair (Enterprise Edition). Consulte Etapa 1: Crie uma instância
Uma lista de permissões de IP configurada para a instância. Consulte Etapa 2: Configure listas de permissões
Limitações
O Realtime Compute for Apache Flink exige o Ververica Runtime (VVR) 6.0.6 ou posterior para usar o conector do Tair (Enterprise Edition).
TairTs, TairCpc, TairRoaring, TairVector e TairGis requerem o VVR 8.0.1 ou posterior.
O conector não suporta a configuração de múltiplos hosts.
Sintaxe
A instrução DDL a seguir cria uma tabela sink do Tair. A definição de PRIMARY KEY é obrigatória.
CREATE TABLE tair_table (
a STRING,
b STRING,
PRIMARY KEY (a) NOT ENFORCED
) WITH (
'connector' = 'tair',
'host' = '<yourHost>',
'mode' = '<dataStructure>'
);
O Tair suporta todas as estruturas de dados do Redis (STRING, LIST, SET, HASHMAP, SORTEDSET), além de suas próprias estruturas proprietárias. Para exemplos de sintaxe das estruturas compatíveis com Redis, consulte Conector do Tair (Redis OSS-compatible).
Opções do conector
Parâmetros obrigatórios
|
Parâmetro |
Tipo |
Descrição |
|
|
String |
Defina como |
|
|
String |
Endpoint do servidor Tair. Use um endpoint interno para evitar limitações de latência e largura de banda das conexões de rede pública. |
|
|
STRING |
Estrutura de dados de destino da gravação. Consulte Estruturas de dados suportadas para verificar os valores válidos. |
Parâmetros de conexão
|
Parâmetro |
Tipo |
Padrão |
Descrição |
|
|
INT |
|
Número da porta do servidor Tair. |
|
|
String |
String vazia |
Senha do banco de dados Tair. Deixe em branco para desativar a verificação por senha. |
|
|
INT |
|
ID do banco de dados de destino. |
|
|
BOOLEAN |
|
Defina como |
Parâmetros de comportamento de gravação
|
Parâmetro |
Tipo |
Padrão |
Descrição |
|
|
BOOLEAN |
|
Controla o comportamento ao receber uma mensagem de retração. |
|
|
STRING |
|
Modo de gravação no sink. |
|
|
STRING |
None |
Quando |
Parâmetros de expiração de chave
|
Parâmetro |
Tipo |
Padrão |
Descrição |
|
|
LONG |
|
Tempo de vida relativo (TTL) das chaves inseridas, em milissegundos. O valor |
|
|
LONG |
|
Timestamp absoluto de expiração para as chaves inseridas, em milissegundos. Só tem efeito quando |
Parâmetros de expiração de campo (apenas TairHash e TairTs)
| Parâmetro | Tipo | Padrão | Descrição |
|---|---|---|---|
fieldExpireMode |
String | None |
Modo de expiração para campos no TairHash ou skeys no TairTs. None: sem expiração. millisecond: expiração relativa definida pelo valor de fieldExpireValue. unixtime: expiração absoluta definida pelo valor de fieldExpireValue. dynamic_millisecond: expiração relativa lida da coluna indicada por fieldExpireValue. dynamic_unixtime: expiração absoluta lida da coluna indicada por fieldExpireValue. Importante
As skeys do TairTs devem usar obrigatoriamente |
fieldExpireValue |
String | None | Se fieldExpireMode for millisecond ou unixtime, este é o valor fixo de expiração. Caso fieldExpireMode seja dynamic_millisecond ou dynamic_unixtime, refere-se ao nome da coluna DDL que contém o valor de expiração. |
Mapeamento de tipos de dados
|
Tipo Flink |
Tipo Tair |
|
VARCHAR |
STRING |
|
DOUBLE |
DOUBLE |
Estruturas de dados suportadas
O parâmetro mode aceita os valores listados abaixo. Cada estrutura de dados possui requisitos específicos de colunas DDL, que variam conforme o incrMode.
Estruturas de dados compatíveis com Redis
Para formatos DDL de STRING, LIST, SET, HASHMAP e SORTEDSET, consulte Conector do Tair (Redis OSS-compatible).
Estruturas de dados proprietárias do Tair
|
Estrutura de dados |
** |
Colunas DDL |
Comando de gravação |
|
|
2: key (STRING), value (STRING) |
|
|
|
TairString |
|
1: key (STRING) |
|
|
TairString |
|
2: key (STRING), incrValue (STRING) |
|
|
|
3: key (STRING), field (STRING), value (STRING) |
|
|
|
TairHash |
|
2: key (STRING), field (STRING) |
|
|
TairHash |
|
3: key (STRING), field (STRING), incrValue (STRING) |
|
|
|
3–258: key (STRING), member (STRING), score×N (DOUBLE, até 256 dimensões) |
|
|
|
TairZset |
|
2: key (STRING), member (STRING) |
|
|
TairZset |
|
3: key (STRING), member (STRING), incrValue (STRING) |
|
|
Deve ser |
2: key (STRING), item (STRING) |
|
|
|
Deve ser |
3: key (STRING), path (STRING), json (STRING) |
|
|
|
|
4: index (STRING), doc_id (STRING), document (STRING, JSON), mapping (STRING) |
|
|
|
TairSearch |
|
4: index (STRING), doc_id (STRING), field (STRING), mapping (STRING) |
|
|
TairSearch |
|
5: index (STRING), doc_id (STRING), field (STRING), mapping (STRING), incrValue (STRING) |
|
|
Deve ser |
2: key (STRING), item (STRING) |
|
|
|
Deve ser |
3: key (STRING), polygon_name (STRING), polygon_wkt (STRING) |
|
|
|
Deve ser |
3: key (STRING), offset (BIGINT), value (BIGINT, 0 ou 1) |
|
|
|
Deve ser |
6: index_name (STRING), pk (STRING), vector_data (STRING), dims (INT), algorithm (STRING), distance_method (STRING) |
|
|
|
|
4: pkey (STRING, grupo de timeline), skey (STRING, timeline única), timestamp (STRING), value (STRING) |
|
|
|
TairTs |
|
3: pkey (STRING), skey (STRING), timestamp (STRING) |
|
|
TairTs |
|
4: pkey (STRING), skey (STRING), timestamp (STRING), incrValue (STRING) |
|
O TairSearch e o TairVector exigem a criação prévia de índice e mapeamento antes da inserção de dados. Para o TairSearch, executeTFT.CREATEINDEX index mappings. Para o TairVector, executeTVS.CREATEINDEX index_name dims algorithm distance_method. O TairBloom cria automaticamente uma chave com capacidade padrão de 100 elementos e taxa de erro de 0,01 na primeira inserção. Na ordenação multidimensional do TairZset, todas as dimensões de score devem usar o mesmo formato.
Exemplos
Gravar dados em uma tabela sink do TairSearch
Este exemplo usa index_name como índice do TairSearch, doc_id como identificador do documento, doc como corpo JSON do documento e mapping como definição de mapeamento do índice.
CREATE TEMPORARY TABLE datagen_stream (
v STRING, -- document content
p STRING -- mapping definition
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
index_name STRING, -- TairSearch index name
doc_id STRING, -- document ID
doc STRING, -- document body (JSON)
mapping STRING, -- index mapping
PRIMARY KEY (index_name) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairsearch',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}'
);
INSERT INTO tair_output
SELECT
'index' AS index_name,
v AS doc_id,
p AS doc,
'{"mappings":{"_source":{"enabled":true},"properties":{"product_id":{"type":"keyword","ignore_above":128},"product_name":{"type":"text"}}}}' AS mapping
FROM datagen_stream;
Gravar dados com incremento dinâmico (TairString com incrMode=dynamic_float)
Este exemplo usa key como chave do TairString e step como valor de incremento por registro, obtido da coluna step.
CREATE TEMPORARY TABLE datagen_stream (
v STRING, -- key
p STRING -- increment value per record
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
key STRING, -- TairString key
step STRING, -- increment value (column referenced by incrValue)
PRIMARY KEY (key) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairstring',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}',
'incrMode' = 'dynamic_float',
'incrValue' = 'step'
);
INSERT INTO tair_output
SELECT *
FROM datagen_stream;
Gravar dados com incremento fixo (TairString com incrMode=float)
Este exemplo aplica um incremento fixo de 11.11 a cada chave em cada operação de gravação.
CREATE TEMPORARY TABLE datagen_stream (
v STRING,
p STRING
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
key STRING, -- TairString key
PRIMARY KEY (key) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairstring',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}',
'incrMode' = 'float',
'incrValue' = '11.11'
);
INSERT INTO tair_output
SELECT v
FROM datagen_stream;
Gravar dados com TTL no nível de campo (TairHash com fieldExpireMode=millisecond)
Este exemplo grava em uma tabela sink do TairHash onde cada campo expira 1000 milissegundos após a gravação.
CREATE TEMPORARY TABLE datagen_stream (
v STRING,
p STRING,
s STRING
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE tair_output (
key STRING, -- TairHash key
field STRING, -- hash field
value STRING, -- field value
PRIMARY KEY (key) NOT ENFORCED
) WITH (
'connector' = 'tair',
'mode' = 'tairhash',
'host' = '${tairHost}',
'port' = '${tairPort}',
'password' = '${password}',
'fieldExpireMode' = 'millisecond',
'fieldExpireValue' = '1000'
);
INSERT INTO tair_output
SELECT v, p, s
FROM datagen_stream;
Próximos passos
Conector do Tair (Redis OSS-compatible) — exemplos de sintaxe para estruturas de dados compatíveis com Redis
Métricas — lista completa de métricas do conector e suas definições