Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Tair (Enterprise)

Última atualização: Jun 27, 2026

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

numBytesSend, numBytesSendPerSecond, numRecordsSend, numRecordsSendPerSecond, numRecordSendErrors, currentSendTime

Para obter detalhes sobre as métricas, consulte Métricas.

Pré-requisitos

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

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

connector

String

Defina como tair.

host

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.

mode

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

port

INT

6379

Número da porta do servidor Tair.

password

String

String vazia

Senha do banco de dados Tair. Deixe em branco para desativar a verificação por senha.

dbNum

INT

0

ID do banco de dados de destino.

clusterMode

BOOLEAN

false

Defina como true para usar a arquitetura de cluster. Defina como false para o modo standalone.

Parâmetros de comportamento de gravação

Parâmetro

Tipo

Padrão

Descrição

ignoreDelete

BOOLEAN

false

Controla o comportamento ao receber uma mensagem de retração. false: exclua os dados inseridos e sua chave. true: mantém os dados inseridos e sua chave.

incrMode

STRING

None

Modo de gravação no sink. None: operação de inserção. int: INCRBY com incremento fixo definido por incrValue. float: INCRBYFLOAT com incremento fixo definido por incrValue. dynamic_int: INCRBY com incremento lido da coluna especificada em incrValue. dynamic_float: INCRBYFLOAT com incremento lido da coluna especificada em incrValue.

incrValue

STRING

None

Quando incrMode é int ou float, representa o valor de incremento fixo. Se incrMode for dynamic_int ou dynamic_float, indica o nome da coluna DDL que contém o valor do incremento. Não usado quando incrMode é None.

Parâmetros de expiração de chave

Parâmetro

Tipo

Padrão

Descrição

expiration

LONG

0

Tempo de vida relativo (TTL) das chaves inseridas, em milissegundos. O valor 0 desativa o TTL.

expirationAt

LONG

0

Timestamp absoluto de expiração para as chaves inseridas, em milissegundos. Só tem efeito quando expiration está definido como 0. O valor 0 desativa a expiração absoluta.

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

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

**incrMode**

Colunas DDL

Comando de gravação

TairString

None

2: key (STRING), value (STRING)

exset key value

TairString

int ou float

1: key (STRING)

exincrby/exincrbyfloat key incrValue

TairString

dynamic_int ou dynamic_float

2: key (STRING), incrValue (STRING)

exincrby/exincrbyfloat key incrValue

TairHash

None

3: key (STRING), field (STRING), value (STRING)

exhset key field value

TairHash

int ou float

2: key (STRING), field (STRING)

exhincrby/exincrbyfloat key field incrValue

TairHash

dynamic_int ou dynamic_float

3: key (STRING), field (STRING), incrValue (STRING)

exhincrby/exincrbyfloat key field incrValue

TairZset

None

3–258: key (STRING), member (STRING), score×N (DOUBLE, até 256 dimensões)

exzadd key score member

TairZset

int ou float

2: key (STRING), member (STRING)

exzincyby key member incrValue

TairZset

dynamic_int ou dynamic_float

3: key (STRING), member (STRING), incrValue (STRING)

exzincyby key member incrValue

TairBloom

Deve ser None

2: key (STRING), item (STRING)

BF.ADD key item

TairDoc

Deve ser None

3: key (STRING), path (STRING), json (STRING)

JSON.SET key path json

TairSearch

None

4: index (STRING), doc_id (STRING), document (STRING, JSON), mapping (STRING)

TFT.ADDDOC index document docid

TairSearch

int ou float

4: index (STRING), doc_id (STRING), field (STRING), mapping (STRING)

TFT.INCRLONGDOCFIELD/TFT.INCRFLOATDOCFIELD index doc_id field increment

TairSearch

dynamic_int ou dynamic_float

5: index (STRING), doc_id (STRING), field (STRING), mapping (STRING), incrValue (STRING)

TFT.INCRLONGDOCFIELD/TFT.INCRFLOATDOCFIELD index doc_id field increment

TairCpc

Deve ser None

2: key (STRING), item (STRING)

CPC.UPDATE key item

TairGis

Deve ser None

3: key (STRING), polygon_name (STRING), polygon_wkt (STRING)

GIS.ADD area polygonName polygonWkt

TairRoaring

Deve ser None

3: key (STRING), offset (BIGINT), value (BIGINT, 0 ou 1)

TR.SETBIT key offset value

TairVector

Deve ser None

6: index_name (STRING), pk (STRING), vector_data (STRING), dims (INT), algorithm (STRING), distance_method (STRING)

TVS.HSET index_name key VECTOR vector_data

TairTs

None

4: pkey (STRING, grupo de timeline), skey (STRING, timeline única), timestamp (STRING), value (STRING)

EXTS.S.RAW_MODIFY Pkey Skey timestamp value

TairTs

float

3: pkey (STRING), skey (STRING), timestamp (STRING)

EXTS.S.RAW_INCRBY Pkey Skey timestamp incrValue

TairTs

dynamic_float

4: pkey (STRING), skey (STRING), timestamp (STRING), incrValue (STRING)

EXTS.S.RAW_INCRBY Pkey Skey timestamp incrValue

O TairSearch e o TairVector exigem a criação prévia de índice e mapeamento antes da inserção de dados. Para o TairSearch, execute TFT.CREATEINDEX index mappings . Para o TairVector, execute TVS.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