Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Tablestore

Última atualização: Jun 27, 2026

O conector do Tablestore permite usar tabelas do Tablestore como tabelas de origem, tabelas de dimensão e tabelas de destino em jobs Flink SQL executados no modo streaming.

Capacidades do conector

Item

Descrição

Modo de execução

Modo streaming

Tipo de API

API SQL

Tipo de tabela

Tabela de origem, tabela de dimensão e tabela de destino

Formato de dados

N/A

Métricas da tabela de destino

numBytesOut, numBytesOutPerSecond, numRecordsOut, numRecordsOutPerSecond, currentSendTime

Atualização ou exclusão de dados na tabela de destino

Compatível

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

Pré-requisitos

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

  • Uma instância do Tablestore e uma tabela do Tablestore. Consulte Usar o Tablestore.

Limites de uso

O acesso entre contas às instâncias do Tablestore é compatível. Ao usar um endpoint de VPC, a instância do Tablestore deve estar na mesma região que o Flink. Defina os parâmetros accessId e accessKey com o par AccessKey da conta proprietária da instância do Tablestore.

Sintaxe

Os três tipos de tabela usam 'connector'='ots' na cláusula WITH, com opções específicas para cada tipo.

Tabela de destino

CREATE TABLE ots_sink (
  name VARCHAR,
  age BIGINT,
  birthday BIGINT,
  PRIMARY KEY (name, age) NOT ENFORCED
) WITH (
  'connector'='ots',
  'instanceName'='<yourInstanceName>',
  'tableName'='<yourTableName>',
  'accessId'='${ak_id}',
  'accessKey'='${ak_secret}',
  'endPoint'='<yourEndpoint>',
  'valueColumns'='birthday'
);
Uma tabela de destino do Tablestore exige uma chave primária. Cada registro de saída é anexado à tabela para atualizar os dados existentes.

Tabela de dimensão

CREATE TABLE ots_dim (
  id INT,
  len INT,
  content STRING
) WITH (
  'connector'='ots',
  'endPoint'='<yourEndpoint>',
  'instanceName'='<yourInstanceName>',
  'tableName'='<yourTableName>',
  'accessId'='${ak_id}',
  'accessKey'='${ak_secret}'
);

Tabela de origem

CREATE TABLE tablestore_stream (
  `order` VARCHAR,
  orderid VARCHAR,
  customerid VARCHAR,
  customername VARCHAR
) WITH (
  'connector'='ots',
  'endPoint'='<yourEndpoint>',
  'instanceName'='flink-source',
  'tableName'='flink_source_table',
  'tunnelName'='flinksourcestream',
  'accessId'='${ak_id}',
  'accessKey'='${ak_secret}',
  'ignoreDelete'='false'
);

Metadados disponíveis

A tabela de origem do Tablestore expõe dois campos de metadados por meio da palavra-chave METADATA. Use esses campos para rastrear o tipo de operação e o momento de cada evento de alteração.

Chave de metadados

Tipo de dados Flink

Descrição

type

STRING

O tipo de operação de dados (mapeado para OtsRecordType).

timestamp

BIGINT

O horário da operação de dados em microssegundos (mapeado para OtsRecordTimestamp). Definido como 0 para leituras completas de dados.

Para ler campos de metadados, declare-os com a sintaxe METADATA FROM:

CREATE TABLE tablestore_stream (
  `order` VARCHAR,
  orderid VARCHAR,
  customerid VARCHAR,
  customername VARCHAR,
  record_type STRING METADATA FROM 'type',
  record_timestamp BIGINT METADATA FROM 'timestamp'
) WITH (
  ...
);

Opções do conector

Opções gerais

Todos os tipos de tabela compartilham as seguintes opções.

Opção

Tipo

Obrigatório

Padrão

Descrição

connector

String

Sim

Defina como ots.

instanceName

String

Sim

Nome da instância do Tablestore.

endPoint

String

Sim

Endpoint da instância do Tablestore. Consulte Endpoints.

tableName

String

Sim

Nome da tabela.

accessId

String

Sim

AccessKey ID da sua conta Alibaba Cloud ou de um usuário do Resource Access Management (RAM). Consulte Como visualizo o AccessKey ID e o AccessKey secret?

accessKey

String

Sim

AccessKey secret da sua conta Alibaba Cloud ou de um usuário RAM.

connectTimeout

Integer

Não

30000

Tempo limite de conexão em milissegundos.

socketTimeout

Integer

Não

30000

Tempo limite do socket em milissegundos.

ioThreadCount

Integer

Não

4

Número de threads de I/O.

callbackThreadPoolSize

Integer

Não

4

Tamanho do pool de threads de callback.

Importante

Use variáveis para armazenar seu par AccessKey em vez de codificá-lo diretamente.

Opções da tabela de origem

Opção

Tipo

Obrigatório

Padrão

Descrição

tunnelName

String

Sim

Nome do túnel do Tablestore. Crie o túnel no console do Tablestore antes de usar esta opção. Tipos de túnel compatíveis: Incremental, Full e Differential. Consulte a seção "Criar um túnel" em Início rápido.

ignoreDelete

Boolean

Não

false

Indica se as operações de exclusão devem ser ignoradas. true: ignorar; false: processar operações de exclusão.

skipInvalidData

Boolean

Não

false

Indica se dados inválidos devem ser ignorados. true: ignorar dados inválidos; false: reportar erro. Requer Ververica Runtime (VVR) 8.0.4 ou posterior.

retryStrategy

Enum

Não

TIME

Política de nova tentativa. TIME: tentar novamente até que retryTimeoutMs expire; COUNT: tentar novamente até atingir retryCount.

retryCount

Integer

Não

3

Número máximo de novas tentativas. Aplica-se quando retryStrategy é COUNT.

retryTimeoutMs

Integer

Não

180000

Tempo limite para novas tentativas em milissegundos. Aplica-se quando retryStrategy é TIME.

streamOriginColumnMapping

String

Não

Mapeamento dos nomes originais das colunas para os nomes reais. Formato: origin_col1:col1,origin_col2:col2.

outputSpecificRowType

Boolean

Não

false

Indica se o tipo de linha específico deve ser repassado. false: todas as linhas são tratadas como INSERT; true: as linhas podem ser INSERT, DELETE ou UPDATE_AFTER.

dataFetchTimeoutMs

Integer

Não

10000

Tempo máximo em milissegundos para buscar dados de uma única partição. Reduza este valor para diminuir a latência geral de sincronização ao sincronizar muitas partições. Requer VVR 8.0.10 ou posterior.

enableRequestCompression

Boolean

Não

false

Indica se a compactação de requisições deve ser ativada. Reduz o uso de largura de banda ao custo de maior carga de CPU. Requer VVR 8.0.10 ou posterior.

Opções da tabela de destino

Opção

Tipo

Obrigatório

Padrão

Descrição

valueColumns

String

Sim

Nomes das colunas a serem gravadas. Separe múltiplos nomes de coluna com vírgulas (,).

retryIntervalMs

Integer

Não

1000

Intervalo entre novas tentativas em milissegundos.

maxRetryTimes

Integer

Não

10

Número máximo de novas tentativas.

bufferSize

Integer

Não

5000

Número máximo de registros armazenados em buffer antes que uma gravação seja acionada.

batchWriteTimeoutMs

Integer

Não

5000

Tempo limite de gravação em milissegundos. Se os registros em buffer não atingirem bufferSize dentro deste período, todos os registros em buffer serão gravados.

batchSize

Integer

Não

100

Número de registros gravados por lote. Máximo: 200.

ignoreDelete

Boolean

Não

false

Indica se as operações de exclusão devem ser ignoradas.

autoIncrementKey

String

Não

Nome da coluna de chave primária com incremento automático. Configure apenas se a tabela de destino possuir uma coluna de chave primária com incremento automático. Requer VVR 8.0.4 ou posterior.

overwriteMode

Enum

Não

PUT

Modo de gravação. PUT: sobrescrever no modo PUT; UPDATE: sobrescrever no modo UPDATE. O modo de coluna dinâmica requer UPDATE.

defaultTimestampInMillisecond

Long

Não

-1

Timestamp padrão para gravações. Se não definido, a hora atual do sistema será utilizada.

dynamicColumnSink

Boolean

Não

false

Indica se o modo de coluna dinâmica deve ser ativado. Neste modo, nenhuma coluna é predefinida; as colunas são inseridas com base em valores de tempo de execução. As primeiras N colunas definem a chave primária. A penúltima coluna contém o nome da coluna e a última coluna contém seu valor — ambas devem ser STRING. Se ativado, overwriteMode deve ser UPDATE e chaves primárias com incremento automático não são compatíveis.

checkSinkTableMeta

Boolean

Não

true

Indica se deve verificar se a chave primária da tabela do Tablestore corresponde à chave primária declarada na instrução CREATE TABLE.

enableRequestCompression

Boolean

Não

false

Indica se a compactação de requisições deve ser ativada durante as gravações.

maxColumnsCount

Integer

Não

128

Número máximo de colunas gravadas na tabela de destino. Se definido acima de 128, ocorre o erro The count of attribute columns exceeds the maximum. Requer VVR 8.0.10 ou posterior.

storageType

String

Não

WIDE_COLUMN

Tipo da tabela de destino. WIDE_COLUMN: tabela de colunas largas; TIMESERIES: tabela de séries temporais.

Opções da tabela de dimensão

Funcionamento do cache

O cache da tabela de dimensão reduz consultas repetidas ao Tablestore. Escolha uma política de cache com base no tamanho da sua tabela e nos padrões de consulta:

  • None: Sem cache. Cada consulta acessa o Tablestore diretamente. Recomendado quando os dados mudam frequentemente e a atualização é crítica.

  • LRU: Armazena em cache um número fixo de registros acessados recentemente. Quando uma consulta não encontra o dado no cache, o conector consulta o Tablestore e atualiza o cache com o resultado. Defina cacheSize e cacheTTLMs ao usar esta política.

  • ALL (padrão): Carrega toda a tabela de dimensão no cache antes do início do job. Todas as consultas subsequentes são atendidas pelo cache. Quando o cache expira (cacheTTLMs), o conector recarrega todos os dados. Use ALL quando a tabela for pequena e você esperar muitas consultas com chaves ausentes. Ao usar ALL, aumente a memória do nó de junção — o cache requer aproximadamente o dobro do tamanho da tabela remota.

Opção

Tipo

Obrigatório

Padrão

Descrição

retryIntervalMs

Integer

Não

1000

Intervalo entre novas tentativas em milissegundos.

maxRetryTimes

Integer

Não

10

Número máximo de novas tentativas.

cache

String

Não

ALL

Política de cache: None, LRU ou ALL.

cacheSize

Integer

Não

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

cacheTTLMs

Integer

Não

TTL do cache em milissegundos. Para LRU: tempo limite por entrada. Para ALL: intervalo de atualização completa do cache. Deixe indefinido para desativar a expiração.

cacheEmpty

Boolean

Não

Indica se resultados vazios (sem correspondência) devem ser armazenados em cache. true: armazenar em cache; false: não armazenar em cache.

cacheReloadTimeBlackList

String

Não

Janelas de tempo durante as quais o cache ALL não é atualizado. Formato: 2017-10-24 14:00 -> 2017-10-24 15:00, 2017-11-10 23:30 -> 2017-11-11 08:00. Separe múltiplas janelas com vírgulas; use -> entre os horários de início e fim.

async

Boolean

Não

false

Indica se a consulta assíncrona deve ser ativada. true: consultas assíncronas (os resultados não são ordenados); false: consultas síncronas.

Mapeamentos de tipos de dados

Tabela de origem

Tipo Tablestore

Tipo Flink SQL

INTEGER

BIGINT

STRING

STRING

BOOLEAN

BOOLEAN

DOUBLE

DOUBLE

BINARY

BINARY

Tabela de destino

Tipo Flink SQL

Tipo Tablestore

BINARY

BINARY

VARBINARY

BINARY

CHAR

STRING

VARCHAR

STRING

TINYINT

INTEGER

SMALLINT

INTEGER

INTEGER

INTEGER

BIGINT

INTEGER

FLOAT

DOUBLE

DOUBLE

DOUBLE

BOOLEAN

BOOLEAN

Exemplos

Ler do Tablestore e gravar no Tablestore

Este exemplo lê dados de pedidos de uma tabela de origem do Tablestore via Tunnel Service e os grava em uma tabela de destino do Tablestore. A tabela de destino usa uma coluna de chave primária com incremento automático.

CREATE TEMPORARY TABLE tablestore_stream (
  `order` VARCHAR,
  orderid VARCHAR,
  customerid VARCHAR,
  customername VARCHAR
) WITH (
  'connector'='ots',
  'endPoint'='<yourEndpoint>',
  'instanceName'='flink-source',
  'tableName'='flink_source_table',
  'tunnelName'='flinksourcestream',
  'accessId'='${ak_id}',
  'accessKey'='${ak_secret}',
  'ignoreDelete'='false',
  'skipInvalidData'='false'
);

CREATE TEMPORARY TABLE ots_sink (
  `order` VARCHAR,
  orderid VARCHAR,
  customerid VARCHAR,
  customername VARCHAR,
  PRIMARY KEY (`order`, orderid) NOT ENFORCED
) WITH (
  'connector'='ots',
  'endPoint'='<yourEndpoint>',
  'instanceName'='flink-sink',
  'tableName'='flink_sink_table',
  'accessId'='${ak_id}',
  'accessKey'='${ak_secret}',
  'valueColumns'='customerid,customername',
  'autoIncrementKey'='${auto_increment_primary_key_name}'
);

INSERT INTO ots_sink
SELECT `order`, orderid, customerid, customername FROM tablestore_stream;

Sincronizar uma tabela de colunas largas para uma tabela de séries temporais

Este exemplo lê de uma tabela de origem de colunas largas e grava em uma tabela de destino de séries temporais. A coluna tags da tabela de destino usa MAP<STRING, STRING> para armazenar pares chave-valor de tags, e storageType está definido como TIMESERIES.

CREATE TEMPORARY TABLE timeseries_source (
  measurement STRING,
  datasource STRING,
  tag_a STRING,
  `time` BIGINT,
  binary_value BINARY,
  bool_value BOOLEAN,
  double_value DOUBLE,
  long_value BIGINT,
  string_value STRING,
  tag_b STRING,
  tag_c STRING,
  tag_d STRING,
  tag_e STRING,
  tag_f STRING
) WITH (
  'connector'='ots',
  'endPoint'='https://iotstore-test.cn-hangzhou.vpc.tablestore.aliyuncs.com',
  'instanceName'='iotstore-test',
  'tableName'='test_ots_timeseries_2',
  'tunnelName'='timeseries_source_tunnel_2',
  'accessId'='${ak_id}',
  'accessKey'='${ak_secret}',
  'ignoreDelete'='true'
);

CREATE TEMPORARY TABLE timeseries_sink (
  measurement STRING,
  datasource STRING,
  tags MAP<STRING, STRING>,
  `time` BIGINT,
  binary_value BINARY,
  bool_value BOOLEAN,
  double_value DOUBLE,
  long_value BIGINT,
  string_value STRING,
  tag_b STRING,
  tag_c STRING,
  tag_d STRING,
  tag_e STRING,
  tag_f STRING,
  PRIMARY KEY (measurement, datasource, tags, `time`) NOT ENFORCED
) WITH (
  'connector'='ots',
  'endPoint'='https://iotstore-test.cn-hangzhou.vpc.tablestore.aliyuncs.com',
  'instanceName'='iotstore-test',
  'tableName'='test_timeseries_sink_table_2',
  'accessId'='${ak_id}',
  'accessKey'='${ak_secret}',
  'storageType'='TIMESERIES'
);

-- Build the tags map from individual tag columns and insert into the time series sink table
INSERT INTO timeseries_sink
SELECT
  measurement,
  datasource,
  MAP['tag_a', tag_a, 'tag_b', tag_b, 'tag_c', tag_c, 'tag_d', tag_d, 'tag_e', tag_e, 'tag_f', tag_f] AS tags,
  `time`,
  binary_value,
  bool_value,
  double_value,
  long_value,
  string_value,
  tag_b,
  tag_c,
  tag_d,
  tag_e,
  tag_f
FROM timeseries_source;

Próximos passos