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 |
|
|
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 |
|
|
STRING |
O tipo de operação de dados (mapeado para |
|
|
BIGINT |
O horário da operação de dados em microssegundos (mapeado para |
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 |
|
|
String |
Sim |
— |
Defina como |
|
|
String |
Sim |
— |
Nome da instância do Tablestore. |
|
|
String |
Sim |
— |
Endpoint da instância do Tablestore. Consulte Endpoints. |
|
|
String |
Sim |
— |
Nome da tabela. |
|
|
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? |
|
|
String |
Sim |
— |
AccessKey secret da sua conta Alibaba Cloud ou de um usuário RAM. |
|
|
Integer |
Não |
30000 |
Tempo limite de conexão em milissegundos. |
|
|
Integer |
Não |
30000 |
Tempo limite do socket em milissegundos. |
|
|
Integer |
Não |
4 |
Número de threads de I/O. |
|
|
Integer |
Não |
4 |
Tamanho do pool de threads de callback. |
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 |
|
|
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. |
|
|
Boolean |
Não |
false |
Indica se as operações de exclusão devem ser ignoradas. |
|
|
Boolean |
Não |
false |
Indica se dados inválidos devem ser ignorados. |
|
|
Enum |
Não |
TIME |
Política de nova tentativa. |
|
|
Integer |
Não |
3 |
Número máximo de novas tentativas. Aplica-se quando |
|
|
Integer |
Não |
180000 |
Tempo limite para novas tentativas em milissegundos. Aplica-se quando |
|
|
String |
Não |
— |
Mapeamento dos nomes originais das colunas para os nomes reais. Formato: |
|
|
Boolean |
Não |
false |
Indica se o tipo de linha específico deve ser repassado. |
|
|
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. |
|
|
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 |
|
|
String |
Sim |
— |
Nomes das colunas a serem gravadas. Separe múltiplos nomes de coluna com vírgulas (,). |
|
|
Integer |
Não |
1000 |
Intervalo entre novas tentativas em milissegundos. |
|
|
Integer |
Não |
10 |
Número máximo de novas tentativas. |
|
|
Integer |
Não |
5000 |
Número máximo de registros armazenados em buffer antes que uma gravação seja acionada. |
|
|
Integer |
Não |
5000 |
Tempo limite de gravação em milissegundos. Se os registros em buffer não atingirem |
|
|
Integer |
Não |
100 |
Número de registros gravados por lote. Máximo: 200. |
|
|
Boolean |
Não |
false |
Indica se as operações de exclusão devem ser ignoradas. |
|
|
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. |
|
|
Enum |
Não |
PUT |
Modo de gravação. |
|
|
Long |
Não |
-1 |
Timestamp padrão para gravações. Se não definido, a hora atual do sistema será utilizada. |
|
|
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, |
|
|
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. |
|
|
Boolean |
Não |
false |
Indica se a compactação de requisições deve ser ativada durante as gravações. |
|
|
Integer |
Não |
128 |
Número máximo de colunas gravadas na tabela de destino. Se definido acima de 128, ocorre o erro |
|
|
String |
Não |
|
Tipo da tabela de destino. |
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
cacheSizeecacheTTLMsao 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 |
|
|
Integer |
Não |
1000 |
Intervalo entre novas tentativas em milissegundos. |
|
|
Integer |
Não |
10 |
Número máximo de novas tentativas. |
|
|
String |
Não |
ALL |
Política de cache: |
|
|
Integer |
Não |
— |
Número máximo de registros em cache. Aplica-se quando |
|
|
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. |
|
|
Boolean |
Não |
— |
Indica se resultados vazios (sem correspondência) devem ser armazenados em cache. |
|
|
String |
Não |
— |
Janelas de tempo durante as quais o cache ALL não é atualizado. Formato: |
|
|
Boolean |
Não |
false |
Indica se a consulta assíncrona deve ser ativada. |
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;