Todos os produtos
Search
Central de documentação

Tablestore:Usar o Realtime Compute for Apache Flink para processar dados do Tablestore

Última atualização: Jun 30, 2026

Este tópico descreve como computar dados do Tablestore com o Realtime Compute for Apache Flink. Tabelas de dados ou de séries temporais do Tablestore podem servir como tabela de source ou de resultados no processamento de dados via Realtime Compute for Apache Flink.

Pré-requisitos

Desenvolver um job de computação em tempo real

Etapa 1: Criar um rascunho SQL

  1. Acesse a página de criação de rascunho.

    1. Faça login no console do Realtime Compute for Apache Flink.

    2. Na coluna Actions do workspace desejado, clique em Console.

    3. No painel de navegação à esquerda, clique em Development > ETL.

  2. Clique em New. Na caixa de diálogo New Draft, selecione Blank Stream Draft e clique em Next.

    Nota

    O Realtime Compute for Apache Flink oferece vários modelos de código e suporta sincronização de dados. Cada modelo atende a cenários específicos e fornece exemplos de código e instruções. Clique em um modelo para conhecer os recursos e a sintaxe relacionada do Realtime Compute for Apache Flink e implementar sua lógica de negócios. Para mais informações, consulte Modelos de código e Modelos de sincronização de dados.

  3. Insira as Job Information.

    Parameter

    Description

    Example

    File Name

    Nome do rascunho a ser criado.

    Nota

    O nome do rascunho deve ser único no projeto atual.

    flink-test

    Location

    Pasta onde o arquivo de código do rascunho será salvo.

    Você também pode clicar no ícone 新建文件夹 à direita de uma pasta existente para criar uma subpasta.

    Draft

    Engine version

    Versão do mecanismo Flink que o rascunho atual utilizará. Para mais informações sobre versões do mecanismo, consulte Notas de versão e Versão do mecanismo.

    vvr-8.0.10-flink-1.17

  4. Clique em Create.

Etapa 2: Escrever o rascunho SQL

Nota

Nesta etapa, o código sincroniza dados de uma tabela de dados para outra. Para mais exemplos de instruções SQL, consulte Exemplos de instruções SQL.

  1. Crie uma tabela temporária para a tabela de source e para a tabela de resultados.

    Para mais informações, consulte Apêndice 1: Conector do Tablestore.

    -- Create a temporary table named tablestore_stream for the source table.
    CREATE TEMPORARY TABLE tablestore_stream(
        `order` VARCHAR,
        orderid VARCHAR,
        customerid VARCHAR,
        customername VARCHAR
    ) WITH (
        'connector' = 'ots', -- Specify the connector type of the source table. The value is ots and cannot be changed. 
        'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com', -- Specify the virtual private cloud (VPC) endpoint of the Tablestore instance. 
        'instanceName' = 'xxx', -- Specify the name of the Tablestore instance. 
        'tableName' = 'flink_source_table', -- Specify the name of the source table. 
        'tunnelName' = 'flink_source_tunnel', -- Specify the name of the tunnel that is created for the source table. 
        'accessId' = 'xxxxxxxxxxx', -- Specify the AccessKey ID of the Alibaba Cloud account or RAM user. 
        'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx', -- Specify the AccessKey secret of the Alibaba Cloud account or RAM user. 
        'ignoreDelete' = 'false' -- Specify whether to ignore the real-time data that is generated by delete operations. In this example, this parameter is set to false. 
    );
    
    -- Create a temporary table named tablestore_sink for the result table.
    CREATE TEMPORARY TABLE tablestore_sink(
       `order` VARCHAR,
        orderid VARCHAR,
        customerid VARCHAR,
        customername VARCHAR,
        PRIMARY KEY (`order`,orderid) NOT ENFORCED -- Specify the primary key. 
    ) WITH (
        'connector' = 'ots', -- Specify the connector type of the result table. The value is ots and cannot be changed. 
        'endPoint'='https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com', -- Specify the VPC endpoint of the Tablestore instance. 
        'instanceName' = 'xxx', -- Specify the name of the Tablestore instance. 
        'tableName' = 'flink_sink_table', -- Specify the name of the result table. 
        'accessId' = 'xxxxxxxxxxx',  -- Specify the AccessKey ID of the Alibaba Cloud account or RAM user. 
        'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx', -- Specify the AccessKey secret of the Alibaba Cloud account or RAM user. 
        'valueColumns'='customerid,customername' --Specify the names of the columns that you want to insert to the result table. 
    );
  2. Escreva a lógica do rascunho.

    A instrução SQL de exemplo a seguir mostra como inserir dados da tabela de source na tabela de resultados:

    -- Insert data from the source table into the result table.
    INSERT INTO tablestore_sink
    SELECT `order`, orderid, customerid, customername FROM tablestore_stream;

Etapa 3: (Opcional) Visualizar informações de configuração

Na aba à direita do editor SQL, visualize as configurações ou defina os parâmetros. A tabela a seguir descreve esses parâmetros.

Tab name

Description

Configurations

  • Engine version: Versão do mecanismo Flink para o rascunho.

  • Additional dependencies: Dependências adicionais necessárias para o job, como UDFs.

    Baixe as dependências do Ververica Runtime (VVR), envie-as na página de arquivos de recursos e selecione os arquivos enviados como additional dependencies. Para mais informações, consulte Apêndice 2: Configurar dependências do VVR.

Structure

  • Data flow diagram: Visualize o fluxo de dados do job.

  • Tree structure diagram: Visualize a linhagem de dados do job.

Versions

Exibe o histórico de versões do rascunho. Para detalhes sobre os recursos na coluna Actions, consulte Gerenciar versões de jobs.

Etapa 4: (Opcional) Executar uma verificação de sintaxe

A validação verifica a semântica SQL do job, a conectividade de rede e os metadados da tabela. Clique em SQL Advice na área de resultados para visualizar possíveis riscos de SQL e sugestões de otimização.

  1. No canto superior direito do editor SQL, clique em Validate.

  2. Na caixa de diálogo Validation, clique em Confirm.

Etapa 5: (Opcional) Depurar o rascunho

Use o recurso de depuração para simular a execução da implantação, verificar saídas e validar a lógica de negócios das instruções SELECT e INSERT. Esse recurso aumenta a eficiência do desenvolvimento e reduz riscos de baixa qualidade de dados.

  1. No canto superior direito do editor SQL, clique em Debug.

  2. Na caixa de diálogo Debug, selecione um cluster de sessão e clique em Next.

    Se nenhum cluster estiver disponível, crie um cluster de sessão. Certifique-se de que o cluster de sessão use a mesma versão de mecanismo do rascunho SQL e esteja em execução. Para mais informações, consulte Criar um cluster de sessão.

  3. Configure os dados de depuração.

    • Se usar dados online, pule esta operação.

    • Para usar dados de depuração, clique em Download Debugging Data Template, preencha o modelo com seus dados e envie o arquivo. Para mais informações, consulte Depurar um job.

  4. Após configure os dados, clique em OK.

Etapa 6: Implantar o rascunho

No canto superior direito do editor SQL, clique em Deploy. Na caixa de diálogo Deploy New Version, configure os parâmetros de implantação e clique em OK.

Nota

Clusters de sessão são adequados para ambientes fora de produção, como desenvolvimento e teste. Implante ou depure rascunhos em um cluster de sessão para melhorar a utilização de recursos do JobManager e acelerar o início da implantação. No entanto, não implante rascunhos destinados ao ambiente de produção em clusters de sessão, pois isso pode causar problemas de estabilidade.

Etapa 7: Iniciar a implantação do rascunho e visualizar o resultado da computação

  1. No painel de navegação à esquerda, clique em O&M > Deployments.

  2. Na coluna Actions da implantação desejada, clique em Start.

    Selecione Start with no state e clique em Start. O status Running indica que a implantação está operando corretamente. Para mais informações sobre parâmetros de inicialização, consulte Iniciar um job.

    Nota
    • Recomendamos configure dois núcleos de CPU e 4 GB de memória para cada TaskManager no Realtime Compute for Apache Flink, maximizando assim a capacidade de computação de cada TaskManager. Um TaskManager consegue escrever 10.000 linhas por segundo.

    • Caso o número de partições na tabela de source seja grande, defina a concorrência para menos de 16 no Realtime Compute for Apache Flink. A taxa de escrita aumenta linearmente conforme a concorrência.

  3. Na página Deployments, visualize o resultado da computação.

    1. Na página O&M > Deployments, clique no nome da implantação desejada.

    2. Na aba Job logs, clique na aba Running task managers e, em seguida, clique na tarefa alvo na coluna Path,ID.

    3. Clique em Logs para visualizar as informações de log.

  4. (Opcional) Cancele uma implantação.

    Ao modifique o código SQL de uma implantação, adicionar ou remover parâmetros da cláusula WITH, ou alterar a versão de uma implantação, implante o rascunho correspondente, cancele a implantação e inicie-a novamente para que as alterações tenham efeito. Se uma implantação falhar e não puder reutilizar os dados de estado para recuperação, ou se você precisar atualize configurações de parâmetros que não entram em vigor dinamicamente, cancele e reinicie a implantação. Para mais informações sobre como cancele uma implantação, consulte Cancelar uma implantação.

Apêndices

Apêndice 1: Conector do Tablestore

O Realtime Compute for Apache Flink fornece um conector do Tablestore integrado para ler, gravar e sincronizar dados do Tablestore.

Tabela de source

Sintaxe DDL

Tabela de dados

O código de exemplo a seguir mostra a instrução DDL para criar uma tabela temporária para a tabela de source:

-- Create a temporary table named tablestore_stream for the source table.
CREATE TEMPORARY TABLE tablestore_stream(
    `order` VARCHAR,
    orderid VARCHAR,
    customerid VARCHAR,
    customername VARCHAR
) WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_source_table',
    'tunnelName' = 'flink_source_tunnel',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'ignoreDelete' = 'false'
);

Tabela de séries temporais

O código de exemplo a seguir mostra a instrução DDL para criar uma tabela temporária para a tabela de source:

-- Create a temporary table named tablestore_stream for the source table.
CREATE TEMPORARY TABLE tablestore_stream(
    _m_name STRING,
    _data_source STRING,
    _tags STRING,
    _time BIGINT,
    binary_value BINARY,
    bool_value BOOLEAN,
    double_value DOUBLE,
    long_value BIGINT,
    string_value STRING
) WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_source_table',
    'tunnelName' = 'flink_source_tunnel',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx'
);

Leia campos de metadados do Tunnel Service, como OtsRecordType e OtsRecordTimestamp, como colunas regulares na tabela de source. A tabela a seguir descreve esses campos.

Field

Mapped field in Realtime Compute for Apache Flink

Description

OtsRecordType

type

Tipo da operação.

OtsRecordTimestamp

timestamp

Hora da operação de dados. Unidade: microssegundos.

Nota

Para que o Realtime Compute for Apache Flink leia todos os dados, defina este campo como 0.

Parâmetro na cláusula WITH

Parameter

Applicable table

Required

Description

connector

Geral

Sim

Tipo de conector da tabela de source. O valor é ots e não pode ser alterado.

endPoint

Geral

Sim

Endpoint da instância do Tablestore. Use um endpoint de VPC. Para mais informações, consulte Endpoints.

instanceName

Geral

Sim

Nome da instância do Tablestore.

tableName

Geral

Sim

Nome da tabela de source no Tablestore.

tunnelName

Geral

Sim

Nome do túnel da tabela de source no Tablestore. Para saber como criar um túnel, consulte Criar um túnel.

accessId

Geral

Sim

Par de AccessKey (AccessKey ID e AccessKey secret) da conta Alibaba Cloud ou do usuário RAM.

Importante

Para proteger seu par de AccessKey, recomendamos usar variáveis para especifique-o. Para mais informações, consulte Gerenciar variáveis.

accessKey

Geral

Sim

connectTimeout

Geral

Não

Tempo limite para o conector do Tablestore se conectar ao Tablestore. Unidade: milissegundos. Valor padrão: 30000.

socketTimeout

Geral

Não

Tempo limite de socket para o conector do Tablestore se conectar ao Tablestore. Unidade: milissegundos. Valor padrão: 30000.

ioThreadCount

Geral

Não

Número de threads de I/O. Valor padrão: 4.

callbackThreadPoolSize

Geral

Não

Tamanho do pool de threads de callback. Valor padrão: 4.

ignoreDelete

Tabela de dados

Não

Define se os dados em tempo real gerados por operações de exclusão devem ser ignorados. Valor padrão: false, indicando que esses dados não são ignorados.

skipInvalidData

Geral

Não

Define se dados inválidos (dirty data) devem ser ignorados. Valor padrão: false, indicando que dados inválidos não são ignorados. Se não forem ignorados, o sistema reportará um erro ao processá-los.

Importante

Somente o Realtime Compute for Apache Flink com VVR 8.0.4 ou posterior suporta este parâmetro.

retryStrategy

Geral

Não

Política de nova tentativa. Valores válidos:

  • TIME: O sistema tenta continuamente até o fim do tempo limite definido pelo parâmetro retryTimeoutMs. Este é o valor padrão.

  • COUNT: O sistema tenta continuamente até atingir o número máximo de tentativas definido pelo parâmetro retryCount.

retryCount

Geral

Não

Número máximo de novas tentativas. Configure este parâmetro se definir retryStrategy como COUNT. Valor padrão: 3.

retryTimeoutMs

Geral

Não

Tempo limite para novas tentativas. Unidade: milissegundos. Configure este parâmetro se definir retryStrategy como TIME. Valor padrão: 180000.

streamOriginColumnMapping

Geral

Não

Mapeamento entre os nomes das colunas na tabela de source e os nomes das colunas na tabela temporária.

Nota

Use dois pontos (:) para separar o nome original da coluna do nome real. Use vírgulas (,) para separar múltiplos mapeamentos. Exemplo: origin_col1:col1,origin_col2:col2.

outputSpecificRowType

Geral

Não

Define se um tipo de linha específico deve ser repassado. Valores válidos:

  • false: não repassa um tipo de linha específico. O tipo de linha de todos os dados é INSERT. Este é o valor padrão.

  • true: repassa um tipo de linha específico. O tipo de linha dos dados pode ser INSERT, DELETE ou UPDATE_AFTER.

Mapeamentos de tipos de dados

Tipo de dados do campo no Tablestore

Tipo de dados do campo no Realtime Compute for Apache Flink

INTEGER

BIGINT

STRING

STRING

BOOLEAN

BOOLEAN

DOUBLE

DOUBLE

BINARY

BINARY

Tabela de resultados

Sintaxe DDL

Tabela de dados

O código de exemplo a seguir mostra a instrução DDL para criar uma tabela temporária para a tabela de resultados:

-- Create a temporary table named tablestore_sink for the result table.
CREATE TEMPORARY TABLE tablestore_sink(
   `order` VARCHAR,
    orderid VARCHAR,
    customerid VARCHAR,
    customername VARCHAR,
    PRIMARY KEY (`order`,orderid) NOT ENFORCED
) WITH (
    'connector' = 'ots',
    'endPoint'='https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_sink_table',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'valueColumns'='customerid,customername'
);
Nota

Especifique o esquema da chave primária e pelo menos uma coluna de atributo para a tabela de resultados do Tablestore. Os dados de saída são anexados à tabela de resultados do Tablestore para atualize os dados da tabela.

Tabela de séries temporais

Uma tabela de resultados de séries temporais requer quatro chaves primárias: _m_name, _data_source, _tags e _time. Especifique essas chaves primárias de três formas: usando parâmetros WITH, usando a definição de chave primária da tabela de resultados ou usando uma chave primária no formato Map. Ao definir a coluna _tags, o método de parâmetro WITH tem a maior prioridade, seguido pelo formato Map e pela definição de chave primária da tabela de resultados.

Usar os parâmetros na cláusula WITH

O código de exemplo a seguir mostra como usar os parâmetros na cláusula WITH para definir a sintaxe DDL:

-- Create a temporary table named tablestore_sink for the result table.
CREATE TEMPORARY TABLE tablestore_sink(
    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,
    PRIMARY KEY(measurement, datasource, tag_a, `time`) NOT ENFORCED
) 
WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_sink_table',
    'accessId' ='xxxxxxxxxxx',
    'accessKey' ='xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'storageType' = 'TIMESERIES',
    'timeseriesSchema' = '{"measurement":"_m_name", "datasource":"_data_source", "tag_a":"_tags", "tag_b":"_tags", "tag_c":"_tags", "tag_d":"_tags", "tag_e":"_tags", "tag_f":"_tags", "time":"_time"}'
);

-- Insert data from the source table into the result table.
INSERT INTO tablestore_sink
    select 
    measurement,
    datasource,
    tag_a,
    `time`,
    binary_value,
    bool_value,
    double_value,
    long_value,
    string_value,
    tag_b,
    tag_c,
    tag_d,
    tag_e,
    tag_f
    from tablestore_stream;

Usar a chave primária no formato Map

O código de exemplo a seguir mostra como usar a chave primária no formato Map para definir a sintaxe DDL:

Nota

O Tablestore fornece o tipo de dados Map do Flink para facilitar a geração da coluna _tags da tabela de séries temporais no modelo TimeSeries. O tipo de dados Map suporta operações de mapeamento, como renomeação de colunas e funções simples. Ao usar Map, certifique-se de que a coluna de chave primária _tags esteja localizada na terceira posição.

-- Create a temporary table named tablestore_sink for the result table.
CREATE TEMPORARY TABLE tablestore_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,
    PRIMARY KEY(measurement, datasource, tags, `time`) NOT ENFORCED
)
WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_sink_table',
    'accessId' ='xxxxxxxxxxx',
    'accessKey' ='xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'storageType' = 'TIMESERIES'
);

-- Insert data from the source table into the result table.
INSERT INTO tablestore_sink
    select 
    measurement,
    datasource,
    MAP[`tag_a`, `tag_b`, `tag_c`, `tag_d`, `tag_e`, `tag_f`] AS tags,
    `time`,
    binary_value,
    bool_value,
    double_value,
    long_value,
    string_value
    from timeseries_source;

Usar a chave primária da tabela SINK

O código de exemplo a seguir mostra como usar a chave primária da tabela SINK para definir a sintaxe DDL. A primeira coluna de chave primária é a coluna _m_name, que especifica o nome da medição. A segunda coluna de chave primária é a coluna _data_source, que especifica a fonte de dados. A última coluna de chave primária é a coluna _time, que especifica o carimbo de data/hora. A coluna de chave primária intermediária é a coluna _tags, que especifica as tags da série temporal.

-- Create a temporary table named tablestore_sink for the result table.
CREATE TEMPORARY TABLE tablestore_sink(
    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
    PRIMARY KEY(measurement, datasource, tag_a, tag_b, tag_c, tag_d, tag_e, tag_f, `time`) NOT ENFORCED
) WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_sink_table',
    'accessId' ='xxxxxxxxxxx',
    'accessKey' ='xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'storageType' = 'TIMESERIES'
);

-- Insert data from the source table into the result table.
INSERT INTO tablestore_sink
    select 
    measurement,
    datasource,
    tag_a,
    tag_b,
    tag_c,
    tag_d,
    tag_e,
    tag_f,
    `time`,
    binary_value,
    bool_value,
    double_value,
    long_value,
    string_value
    from timeseries_source;
Parâmetro na cláusula WITH

Parameter

Applicable table

Required

Description

connector

Geral

Sim

Tipo de conector da tabela de resultados. O valor é ots e não pode ser alterado.

endPoint

Geral

Sim

Endpoint da instância do Tablestore. Use um endpoint de VPC. Para mais informações, consulte Endpoints.

instanceName

Geral

Sim

Nome da instância do Tablestore.

tableName

Geral

Sim

Nome da tabela de séries temporais no Tablestore.

accessId

Geral

Sim

Par de AccessKey (AccessKey ID e AccessKey secret) da conta Alibaba Cloud ou do usuário RAM.

Importante

Para proteger seu par de AccessKey, recomendamos usar variáveis para especifique-o. Para mais informações, consulte Gerenciar variáveis.

accessKey

Geral

Sim

valueColumns

Tabela de dados

Sim

Nomes das colunas nas quais os dados serão gravados. Separe vários nomes de coluna com vírgulas (,). Exemplo: ID,NAME.

storageType

Geral

Não

Importante

Se usar uma tabela de séries temporais como tabela de resultados, defina este parâmetro como TIMESERIES.

Tipo da tabela. Valores válidos:

  • WIDE_COLUMN: tabela de dados. Este é o valor padrão.

  • TIMESERIES: tabela de séries temporais.

timeseriesSchema

Tabela de séries temporais

Não

Importante

Ao usar uma tabela de séries temporais como tabela de resultados, se você usar os parâmetros na cláusula WITH para especifique a chave primária da tabela de séries temporais, deverá configure este parâmetro.

Colunas que você deseja especifique como colunas de chave primária da tabela de séries temporais.

  • Especifique a chave primária da tabela de séries temporais como pares chave-valor no formato JSON. Exemplo: {"measurement":"_m_name","datasource":"_data_source","tag_a":"_tags","tag_b":"_tags","tag_c":"_tags","tag_d":"_tags","tag_e":"_tags","tag_f":"_tags","time":"_time"}.

  • Os tipos das colunas de chave primária especificadas devem ser iguais aos tipos das colunas de chave primária na tabela de séries temporais. A coluna de chave primária tags pode consistir em várias colunas.

connectTimeout

Geral

Não

Tempo limite para o conector do Tablestore se conectar ao Tablestore. Unidade: milissegundos. Valor padrão: 30000.

socketTimeout

Geral

Não

Tempo limite de socket para o conector do Tablestore se conectar ao Tablestore. Unidade: milissegundos. Valor padrão: 30000.

ioThreadCount

Geral

Não

Número de threads de I/O. Valor padrão: 4.

callbackThreadPoolSize

Geral

Não

Tamanho do pool de threads de callback. Valor padrão: 4.

retryIntervalMs

Geral

Não

Intervalo entre novas tentativas. Unidade: milissegundos. Valor padrão: 1000.

maxRetryTimes

Geral

Não

Número máximo de novas tentativas. Valor padrão: 10.

bufferSize

Geral

Não

Número máximo de registros de dados que podem ser armazenados no buffer antes que os dados sejam gravados na tabela de resultados. Valor padrão: 5000, indicando que os dados são gravados na tabela de resultados quando o número de registros no buffer atinge 5.000.

batchWriteTimeoutMs

Geral

Não

Tempo limite de gravação. Unidade: milissegundos. Valor padrão: 5000, indicando que todos os dados no buffer são gravados na tabela de resultados se o número de registros no buffer não atingir o valor especificado pelo parâmetro bufferSize dentro de 5.000 milissegundos.

batchSize

Geral

Não

Número de registros de dados que podem ser gravados na tabela de resultados simultaneamente. Valor padrão: 100. Valor máximo: 200.

ignoreDelete

Geral

Não

Define se os dados em tempo real gerados por operações de exclusão devem ser ignorados. Valor padrão: false, indicando que esses dados não são ignorados.

Importante

Se usar uma tabela de dados como tabela de source, configure este parâmetro conforme suas necessidades de negócios.

autoIncrementKey

Tabela de dados

Não

Nome da coluna de chave primária de incremento automático da tabela de resultados, caso ela contenha tal coluna. Se a tabela de resultados não tiver uma coluna de chave primária de incremento automático, não é necessário configure este parâmetro.

Importante

Somente o Realtime Compute for Apache Flink com VVR 8.0.4 ou posterior suporta este parâmetro.

overwriteMode

Geral

Não

Modo de sobrescrita de dados. Valores válidos:

  • PUT: Os dados são gravados na tabela do Tablestore no modo PUT. Este é o valor padrão.

  • UPDATE: Os dados são gravados na tabela do Tablestore no modo UPDATE.

Nota

Somente o modo UPDATE é suportado no modo de coluna dinâmica.

defaultTimestampInMillisecond

Geral

Não

Carimbo de data/hora padrão usado para gravar dados na tabela do Tablestore. Se deixar este parâmetro vazio, o carimbo de data/hora atual do sistema será usado.

dynamicColumnSink

Geral

Não

Define se o modo de coluna dinâmica deve ser ativado. Valor padrão: false, indicando que o modo de coluna dinâmica está desativado.

Nota
  • O modo de coluna dinâmica é adequado para cenários em que nenhuma coluna é especificada para uma tabela e as colunas de dados são inseridas na tabela com base no status da implantação. Especifique as primeiras colunas como colunas de chave primária na instrução de criação da tabela. O valor da penúltima coluna é usado como variável de nome de coluna, o valor da última coluna é usado como valor dessa variável, e o tipo de dados da penúltima coluna deve ser String.

  • Se ative o modo de coluna dinâmica, o recurso de coluna de chave primária de incremento automático não será suportado e você deverá defina o parâmetro overwriteMode como UPDATE.

checkSinkTableMeta

Geral

Não

Define se os metadados da tabela de resultados devem ser verifique. Valor padrão: true, indicando que o sistema verifica se as colunas de chave primária da tabela do Tablestore são iguais às colunas de chave primária especificadas na instrução de criação da tabela.

enableRequestCompression

Geral

Não

Define se a compactação de dados deve ser ativado durante a gravação. Valor padrão: false, indicando que a compactação de dados está desativada durante a gravação.

Mapeamentos de tipos de dados

Tipo de dados do campo no Realtime Compute for Apache Flink

Tipo de dados do campo no Tablestore

BINARY

BINARY

VARBINARY

BINARY

CHAR

STRING

VARCHAR

STRING

TINYINT

INTEGER

SMALLINT

INTEGER

INTEGER

INTEGER

BIGINT

INTEGER

FLOAT

DOUBLE

DOUBLE

DOUBLE

BOOLEAN

BOOLEAN

Exemplos de instruções SQL

Sincronizar dados da tabela de source para a tabela de resultados

Sincronizar dados de uma tabela de dados para uma tabela de séries temporais

Leia dados de uma tabela de dados chamada flink_source_table e grave-os em uma tabela de séries temporais chamada flink_sink_table.

Instrução SQL de exemplo:

-- Create a temporary table named tablestore_stream for the source table.
CREATE TEMPORARY TABLE tablestore_stream(
    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://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_source_table',
    'tunnelName' = 'flink_source_tunnel',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'ignoreDelete' = 'false' 
);

-- Create a temporary table named tablestore_sink for the result table by using the parameters in the WITH clause.
CREATE TEMPORARY TABLE tablestore_sink(
     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,
     PRIMARY KEY(measurement, datasource, tag_a, `time`) NOT ENFORCED
 ) WITH (
     'connector' = 'ots',
     'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
     'instanceName' = 'xxx',
     'tableName' = 'flink_sink_table',
     'accessId' = 'xxxxxxxxxxx',
     'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
     'storageType' = 'TIMESERIES',
     'timeseriesSchema' = '{"measurement":"_m_name","datasource":"_data_source","tag_a":"_tags","tag_b":"_tags","tag_c":"_tags","tag_d":"_tags","tag_e":"_tags","tag_f":"_tags","time":"_time"}'
 );
 
-- Insert data from the source table into the result table.
INSERT INTO tablestore_sink
    select 
    measurement,
    datasource,
    tag_a,
    `time`,
    binary_value,
    bool_value,
    double_value,
    long_value,
    string_value,
    tag_b,
    tag_c,
    tag_d,
    tag_e,
    tag_f
    from tablestore_stream;

Sincronizar dados de uma tabela de séries temporais para uma tabela de dados

Leia dados de uma tabela de séries temporais chamada flink_source_table e grave-os em uma tabela de dados chamada flink_sink_table.

Instrução SQL de exemplo:

-- Create a temporary table named tablestore_stream for the source table.
CREATE TEMPORARY TABLE tablestore_stream(
    _m_name STRING,
    _data_source STRING,
    _tags STRING,
    _time BIGINT,
    binary_value BINARY,
    bool_value BOOLEAN,
    double_value DOUBLE,
    long_value BIGINT,
    string_value STRING
) WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_source_table',
    'tunnelName' = 'flink_source_tunnel',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx'
);

-- Create a temporary table named print_table for the result table. 
CREATE TEMPORARY TABLE tablestore_target(
    measurement STRING,
    datasource STRING,
    tags STRING,
    `time` BIGINT,
    binary_value BINARY,
    bool_value BOOLEAN,
    double_value DOUBLE,
    long_value BIGINT,
    string_value STRING,
    PRIMARY KEY (measurement,datasource, tags, `time`) NOT ENFORCED
) WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_sink_table',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'valueColumns'='binary_value,bool_value,double_value,long_value,string_value'
);

-- Insert data from the source table into the result table.
INSERT INTO tablestore_target
SELECT
    _m_name,
    _data_source,
    _tags,
    _time,
    binary_value,
    bool_value,
    double_value,
    long_value,
    string_value
    from tablestore_stream;

Ler dados da tabela de source e imprimi-los no console do Tablestore

Leia dados da tabela de source chamada flink_source_table em lotes. Use o recurso de depuração de implantação para simular a execução de uma implantação. O resultado da depuração é exibido na parte inferior do editor SQL.

Instrução SQL de exemplo:

-- Create a temporary table named tablestore_stream for the source data table.
CREATE TEMPORARY TABLE tablestore_stream(
    `order` VARCHAR,
    orderid VARCHAR,
    customerid VARCHAR,
    customername VARCHAR
) WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_source_table',
    'tunnelName' = 'flink_source_tunnel',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'ignoreDelete' = 'false' 
);

-- Read data from the source table.
SELECT * FROM tablestore_stream LIMIT 100;

Ler dados da tabela de source e imprimi-los no log do TaskManager

Leia dados da tabela de source chamada flink_source_table e imprima os resultados no log do TaskManager usando o conector Print.

Instrução SQL de exemplo:

-- Create a temporary table named tablestore_stream for the source data table.
CREATE TEMPORARY TABLE tablestore_stream(
    `order` VARCHAR,
    orderid VARCHAR,
    customerid VARCHAR,
    customername VARCHAR
) WITH (
    'connector' = 'ots',
    'endPoint' = 'https://xxx.cn-hangzhou.vpc.tablestore.aliyuncs.com',
    'instanceName' = 'xxx',
    'tableName' = 'flink_source_table',
    'tunnelName' = 'flink_source_tunnel',
    'accessId' = 'xxxxxxxxxxx',
    'accessKey' = 'xxxxxxxxxxxxxxxxxxxxxxxxxxxx',
    'ignoreDelete' = 'false' 
);

-- Create a temporary table named print_table for the result table. 
CREATE TEMPORARY TABLE print_table(
   `order` VARCHAR,
    orderid VARCHAR,
    customerid VARCHAR,
    customername VARCHAR
) WITH (
  'connector' = 'print',   -- Use the Print connector.
  'logger' = 'true'        -- Display the computing result in the development console of Realtime Compute for Apache Flink.
);

-- Print the fields of the source table.
INSERT INTO print_table
SELECT `order`,orderid,customerid,customername from tablestore_stream;

Apêndice 2: Configure dependências do VVR

  1. Baixe as dependências do VVR.

  2. Envie as dependências do VVR.

    1. Faça login no console do Realtime Compute for Apache Flink.

    2. Localize o workspace desejado e clique em Console na coluna Actions.

    3. No painel de navegação à esquerda, clique em Artifacts.

    4. Na página Artifacts, clique em Upload Artifact e selecione o pacote JAR de dependência do VVR.

  3. No lado direito do editor SQL do job desejado, clique na aba Configurations. No campo Additional Dependencies, selecione o pacote JAR de dependência do VVR.