Um nó de lote Flink SQL permite usar instruções SQL padrão para definir e executar tarefas de processamento de dados. Utilize-o para analisar e transformar grandes conjuntos de dados em atividades como limpeza e agregação. O nó oferece configuração visual e representa uma solução eficiente e flexível para processamento em lote de grande escala. Este tópico descreve como utilizar um nó de lote Flink SQL para processar dados em lotes.
Pré-requisitos
Você criou um workspace e vinculou um recurso de computação do Realtime Compute for Apache Flink na seção Administration. Para mais informações, consulte Bind computing resources.
Você criou um nó de lote Flink SQL. Para mais informações, consulte Create a node for a scheduling workflow.
-
Você concedeu as seguintes permissões de API ao usuário RAM ou à função RAM que o DataWorks utiliza para chamar as APIs do Realtime Compute for Apache Flink. Essas permissões são necessárias para enviar e implantar tarefas de nó em um cluster Flink. Para mais informações, consulte Grant permissions.
{ "Version": "1", "Statement": [ { "Effect": "Allow", "Action": ["stream:CreateDeployment", "stream:UpdateDeployment", "stream:GetDeployment", "stream:DeleteDeployment"], "Resource": ["*"] } ] }
Limitação
Somente grupos de recursos serverless são compatíveis. O grupo de recursos exclusivo legado para agendamento não é suportado.
Etapa 1: Desenvolver o nó de lote Flink SQL
Na página de edição do nó de lote Flink SQL, desenvolva a tarefa do nó.
Desenvolver código SQL
Desenvolva o código da tarefa na área de edição SQL. No código, defina variáveis usando o formato ${variable_name}. Em seguida, no lado direito da página de edição do nó, atribua um valor à variável na seção Scheduling Parameters do painel Scheduling Settings. Isso permite passar parâmetros dinamicamente para o código em cenários de agendamento. Para mais informações sobre o uso de parâmetros de agendamento, consulte Scheduling parameter sources and their expressions. Veja um exemplo abaixo.
-- Create a source table named datagen_source.
CREATE TEMPORARY TABLE datagen_source_${var}(
name VARCHAR
) WITH (
'connector' = 'datagen',
'number-of-rows' = '1000'
);
-- Create a result table named blackhole_sink.
CREATE TEMPORARY TABLE blackhole_sink_${var}(
name VARCHAR
) WITH (
'connector' = 'blackhole'
);
-- Insert data from the source table into the result table.
INSERT INTO blackhole_sink_${var}
SELECT
name
FROM datagen_source_${var};
Neste exemplo, o parâmetro bizdate tem o valor $[yyyymmdd], o que possibilita a sincronização em lote de novos dados diários.
Etapa 2: Configurar o nó de lote Flink SQL
Configure os parâmetros da tarefa para o nó de lote Flink SQL conforme suas necessidades de negócio.
Configurar recursos Flink
Defina os seguintes parâmetros no lado direito da página de edição, na seção Flink resource information dentro de Scheduling Settings. Para mais detalhes, consulte Configure schedule settings.
|
Parâmetro |
Descrição |
|
Flink cluster |
Nome do recurso de computação Flink totalmente gerenciado vinculado em Administration. |
|
Flink engine version |
Selecione uma versão do motor adequada às suas necessidades. |
|
Resource Group for Scheduling |
Selecione um serverless resource group com conectividade de rede ao Flink. |
|
Job Manager CPU |
Seguindo as melhores práticas do Flink, o JobManager precisa de pelo menos 0,5 núcleo de CPU e 2 GiB de memória para operar com estabilidade. Recomenda-se 1 núcleo de CPU e 4 GiB de memória, com limite máximo de 16 núcleos de CPU. Ajuste a configuração conforme a escala do cluster e a complexidade do job. |
|
Job Manager Memory |
A configuração de memória do JobManager influencia sua capacidade de lidar com tarefas de agendamento e gerenciamento. Para garantir operação estável e eficiente, utilize valores entre 2 GiB e 64 GiB. Adapte o valor segundo a escala do cluster e os requisitos do job. |
|
Task Manager CPU |
A configuração de CPU do TaskManager determina sua capacidade de processamento de tarefas. De acordo com as melhores práticas do Flink, recomenda-se no mínimo 0,5 núcleo de CPU e 2 GiB de memória, sendo ideal 1 núcleo de CPU e 4 GiB de memória, até o limite de 16 núcleos. Ajuste conforme necessário. |
|
Task Manager Memory |
A memória configurada no TaskManager define o volume de dados e o desempenho no processamento de tarefas. Para assegurar execução estável e eficiente, a memória deve ser de no mínimo 2 GiB, podendo chegar a 64 GiB. |
|
Concurrency |
Define o número de execuções paralelas de tarefas em um job Flink. Maior concorrência pode aumentar a velocidade de processamento e a utilização de recursos. Defina este valor considerando os recursos do cluster e as características do job. |
|
Maximum number of slots |
Um slot é uma unidade de recurso de tamanho fixo em um Task Manager, alocável para tarefas. Cada slot executa uma instância de tarefa ou operador. Ajuste o número máximo de slots conforme os recursos disponíveis. |
|
Number of slots per TaskManager |
A quantidade de slots por TaskManager indica quantas tarefas ele processa simultaneamente. Modifique essa configuração para otimizar o uso de recursos e a capacidade de processamento paralelo. |
(Opcional) Configurar parâmetros de agendamento
No lado direito da página de edição, na seção Scheduling Parameters dentro de Scheduling Settings, clique em Add parameters e edite o Parameter name e o Parameter Value para usá-los dinamicamente no código.
(Opcional) Configurar parâmetros de runtime do Flink
Configure parâmetros de runtime no lado direito da página de edição, na seção Flink running parameters dentro de Scheduling Settings. Para mais informações, consulte Configure schedule settings.
Ao configurar parâmetros de runtime do Flink, a sintaxe é compatível com VVP (Ververica Platform). Escreva as configurações diretamente em formato YAML, sem adicionar ponto e vírgula ou outros caracteres especiais para quebras de linha.
Para executar a tarefa do nó em um cronograma periódico, configure as informações de agendamento (Scheduling Policy, Scheduling time, Scheduling Dependency e Node output parameters) conforme suas necessidades de negócio. Para mais informações, consulte Configure schedule settings.
Após concluir a configuração da tarefa, clique em Save.
Etapa 3: (Opcional) Depurar o nó de lote Flink SQL
Antes de implantar o nó no ambiente de produção, utilize o recurso de depuração para realizar uma execução de teste do código com dados simulados enviados. Isso permite verificar a lógica SQL e os dados upstream e downstream sem implantar a tarefa no Operation Center.
O recurso de depuração está disponível mediante lista de permissões. Para utilizá-lo, envie um ticket solicitando acesso.
Configurar informações de recursos Flink
No lado direito da página de edição do nó, na seção Flink resource information do painel Run Configuration, configure os parâmetros descritos na tabela a seguir.
|
Parâmetro |
Descrição |
|
Flink Debug Cluster |
Cluster de sessão Flink usado para executar a tarefa de depuração. Este parâmetro é obrigatório. A lista suspensa exibe os clusters de sessão existentes sob o recurso de computação atual, juntamente com seu status de execução. Apenas clusters com status Running podem ser selecionados. Caso não haja clusters disponíveis na lista, clique em Create Cluster para acessar o console do Realtime Compute for Apache Flink e criar um cluster de sessão. |
|
Flink Engine Version |
Versão do motor Flink do cluster de sessão selecionado. O sistema exibe esse valor automaticamente com base no cluster, sem necessidade de entrada manual. |
|
Timeout |
Duração máxima de uma única tarefa de depuração, em minutos. O valor padrão é 30 minutos. A tarefa de depuração é interrompida automaticamente após o tempo especificado. |
Ao alternar o recurso de computação do nó atual, o Flink Debug Cluster selecionado e os dados de depuração enviados são limpos. É necessário selecionar novamente um cluster e reenviar os dados.
Preparar dados de depuração
Na seção Debug Data do painel Run Configuration, prepare dados simulados para as tabelas de origem referenciadas no código.
Clique em Generate Template. O sistema analisa as tabelas de origem referenciadas no código SQL atual e gera registros correspondentes aos nomes das tabelas na lista abaixo. Dados enviados anteriormente não são apagados.
Na coluna Actions do registro da tabela, clique em Download Template para baixar um modelo CSV compatível com a estrutura da tabela de origem.
Preencha os dados de depuração no modelo baixado respeitando a ordem das colunas e salve o arquivo em formato CSV.
Na coluna Actions do registro da tabela, clique em Upload e selecione o arquivo CSV preenchido para envio. Após o sucesso do upload, a coluna Status exibirá Enabled.
(Opcional) Após o sucesso do upload, clique em Preview no painel inferior para visualizar os dados. Para modificar os dados, envie novamente um arquivo CSV para substituir os existentes.
Se não desejar usar os dados simulados de uma tabela de origem específica na sessão de depuração atual, clique em Disable para alterar o status para Disabled. Para reativar os dados, clique em Enable. Somente dados com status Enabled são utilizados na sessão de depuração atual.
Antes de enviar dados de depuração, selecione primeiro um Flink Debug Cluster. Caso contrário, o sistema solicitará a seleção de um recurso de computação.
Os dados de depuração aceitam apenas o formato CSV, e o tamanho do arquivo não pode exceder 1 MB. A primeira linha do arquivo CSV deve conter os nomes das colunas. Recomendamos o uso da codificação UTF-8.
Executar a tarefa de depuração
Com os dados de depuração prontos, clique no botão Run na barra de ferramentas do editor (ou pressione F8). O sistema envia o código, os dados simulados e as informações de recursos Flink ao cluster de sessão selecionado para execução.
Se o código utilizar parâmetros no formato ${variable_name}, certifique-se de ter atribuído valores às variáveis na seção Scheduling Parameters. Durante a depuração, o sistema substitui os placeholders no código pelos valores atribuídos antes do envio.
Visualizar os resultados da depuração
Após a execução da tarefa de depuração, a área de resultados na parte inferior do nó fornece as seguintes informações para ajudar na identificação rápida de problemas:
Code: O código SQL enviado ao motor Flink nesta execução (com as substituições de variáveis aplicadas).
Logs: O log de runtime e informações de erro da tarefa de depuração.
Query results: Os dados de saída da tarefa de depuração.
Etapa 4: Implantar e gerenciar o nó de lote Flink SQL
Após configurar a tarefa do nó, implante-o. Para mais informações, consulte Deploy a node.
Depois de implantar a tarefa, clique em Go to operation and maintenance abaixo de Deploy to Production para visualizar o status de execução das tarefas agendadas no Operation Center. Para mais informações, consulte View scheduled tasks.