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
O Tablestore está ativado e há uma instância criada. Para mais informações, consulte Ativar o Tablestore e criar uma instância.
Crie uma tabela de source, uma tabela de resultados e um túnel para a tabela de source. Para mais informações, consulte Operações em uma tabela de dados, Operações em tabelas de séries temporais e Criar um túnel.
-
Crie um workspace do Realtime Compute for Apache Flink. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
ImportanteO workspace do Realtime Compute for Apache Flink e a instância do Tablestore devem residir na mesma região. Para obter informações sobre as regiões compatíveis com o Realtime Compute for Apache Flink, consulte Regiões.
-
Obtenha um par de AccessKey.
ImportantePor motivos de segurança, recomendamos usar os recursos do Tablestore como usuário do Resource Access Management (RAM). Para mais informações, consulte Usar o par de AccessKey de um usuário RAM para acessar o Tablestore.
Desenvolver um job de computação em tempo real
Etapa 1: Criar um rascunho SQL
-
Acesse a página de criação de rascunho.
Faça login no console do Realtime Compute for Apache Flink.
Na coluna Actions do workspace desejado, clique em Console.
No painel de navegação à esquerda, clique em Development > ETL.
-
Clique em New. Na caixa de diálogo New Draft, selecione Blank Stream Draft e clique em Next.
NotaO 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.
-
Insira as Job Information.
Parameter
Description
Example
File Name
Nome do rascunho a ser criado.
NotaO 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
Clique em Create.
Etapa 2: Escrever o rascunho SQL
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.
-
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. ); -
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 |
|
Structure |
|
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.
No canto superior direito do editor SQL, clique em Validate.
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.
No canto superior direito do editor SQL, clique em Debug.
-
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.
-
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.
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.
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
No painel de navegação à esquerda, clique em O&M > Deployments.
-
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.
NotaRecomendamos 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.
-
Na página Deployments, visualize o resultado da computação.
Na página O&M > Deployments, clique no nome da implantação desejada.
Na aba Job logs, clique na aba Running task managers e, em seguida, clique na tarefa alvo na coluna Path,ID.
Clique em Logs para visualizar as informações de log.
-
(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.