Este tutorial apresenta um exemplo prático de uso dos principais recursos do Realtime Compute for Apache Flink para processamento em lote.
Recursos
O Realtime Compute for Apache Flink oferece os seguintes recursos essenciais para o processamento em lote:
Mapa de desenvolvimento de jobs: Na aba Drafts da página SQL Editor, crie um rascunho em lote e implante-o e execute-o como uma implantação em lote.
Gerenciamento de implantações: Na página Deployments, implante diretamente jobs JAR ou Python em lote. Selecione BATCH na lista suspensa superior para visualizar as implantações em lote ativas. Expanda a implantação desejada para ver a lista de instâncias de job. Geralmente, diferentes instâncias de uma mesma implantação compartilham a lógica de processamento, mas utilizam parâmetros distintos, como a data dos dados processados.
Consulta de dados: Na aba Scripts da página SQL Editor, execute instruções DDL ou consultas curtas para gerenciar e explorar dados rapidamente. Essas consultas rodam em uma sessão Flink pré-criada, permitindo executar consultas simples com baixa latência ao reaproveitar recursos.
Gerenciamento de dados: Na página Catalogs, crie e visualize catálogos, incluindo bancos de dados e tabelas. Para aumentar a eficiência do desenvolvimento, acesse esses itens também na aba Catalogs da página SQL Editor.
Orquestração de tarefas (Public Preview): Na página Workflows, utilize uma interface visual para definir fluxos de trabalho que orquestram dependências entre implantações em lote. Acione o fluxo manualmente ou agende-o para execução periódica.
Gerenciar filas de recursos: Na página Queue Management, particione os recursos dentro de um workspace para evitar contenção entre implantações de streaming e em lote, bem como entre implantações com prioridades diferentes.
Pré-requisitos
Você criou um workspace do Realtime Compute for Apache Flink. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
Você ativou o Object Storage Service (OSS). Para mais detalhes, consulte Início rápido do console. A classe de armazenamento do bucket OSS deve ser Standard. Consulte Classes de armazenamento para saber mais.
Este tutorial utiliza o Apache Paimon para armazenar dados e exige o Ververica Runtime (VVR) 8.0.5 ou superior.
Cenário de exemplo
Este tutorial simula uma plataforma de e-commerce cujos dados estão armazenados no formato lakehouse do Apache Paimon. A estrutura de data warehouse possui múltiplas camadas: Operational Data Store (ODS), Data Warehouse Detail (DWD) e Data Warehouse Service (DWS). Utilize as capacidades de processamento em lote do Flink para processar e limpar os dados antes de gravá-los nas tabelas Paimon e construir uma estrutura de dados em camadas.

Preparações
-
Crie um script.
Na aba Scripts, crie um catálogo contendo bancos de dados e tabelas e insira dados de amostra nessas tabelas.
-
Crie um catálogo Paimon.
-
No editor da aba Scripts, insira a seguinte instrução SQL.
CREATE CATALOG `my_catalog` WITH ( 'type' = 'paimon', 'metastore' = 'filesystem', 'warehouse' = '<warehouse>', 'fs.oss.endpoint' = '<fs.oss.endpoint>', 'fs.oss.accessKeyId' = '<fs.oss.accessKeyId>', 'fs.oss.accessKeySecret' = '<fs.oss.accessKeySecret>' );A tabela a seguir descreve os parâmetros.
Parâmetro
Descrição
Obrigatório
Observações
type
Tipo do catálogo.
Sim
Defina o valor como
paimon.metastore
Tipo de metastore.
Sim
Neste tutorial, defina o valor como
filesystem. Para mais informações sobre outros tipos, consulte Gerenciar catálogos Paimon.warehouse
Diretório do data warehouse no Object Storage Service (OSS).
Sim
O formato é
oss://<bucket>/<object>, onde:bucket: Nome do seu bucket OSS.object: Caminho onde seus dados estão armazenados.
Encontre o bucket e o caminho do objeto no console do OSS.
fs.oss.endpoint
Endpoint do Object Storage Service (OSS).
Não
Obrigatório se o bucket OSS especificado em
warehouseestiver em uma região diferente do seu workspace Flink ou pertencer a outra conta Alibaba Cloud.Para mais informações, consulte Regiões e endpoints.
fs.oss.accessKeyId
AccessKey ID de uma conta Alibaba Cloud ou usuário RAM com permissões de leitura e gravação no bucket OSS.
Não
Obrigatório se o bucket OSS especificado em
warehouseestiver em uma região diferente do seu workspace Flink ou pertencer a outra conta Alibaba Cloud. Para obter um par AccessKey, consulte Criar um par AccessKey.fs.oss.accessKeySecret
AccessKey secret correspondente ao AccessKey ID.
Não
-
Selecione o código e clique em Execute à esquerda.
A mensagem
The following statement has been executed successfully!indica que o catálogo foi criado com sucesso. Visualize-o na página Catalogs ou na aba Catalogs da página SQL Editor.
-
Procedimento
Etapa 1: Criar tabelas ODS e inserir dados de teste
Para simplificar este tutorial, inserimos dados de teste diretamente nas tabelas ODS para popular as tabelas DWD e DWS subsequentes. Em ambientes de produção, geralmente se usa o processamento de streaming do Flink para ler dados de fontes externas e gravá-los no data lake como camada ODS. Para mais detalhes, consulte Introdução ao Paimon: Recursos básicos.
-
No editor da aba Scripts, insira as seguintes instruções SQL e clique em Run à esquerda.
CREATE DATABASE `my_catalog`.`order_dw`; USE `my_catalog`.`order_dw`; CREATE TABLE orders ( order_id BIGINT, user_id STRING, shop_id BIGINT, product_id BIGINT, buy_fee BIGINT, create_time TIMESTAMP, update_time TIMESTAMP, state INT ); CREATE TABLE orders_pay ( pay_id BIGINT, order_id BIGINT, pay_platform INT, create_time TIMESTAMP ); CREATE TABLE product_catalog ( product_id BIGINT, catalog_name STRING ); -- Insert test data INSERT INTO orders VALUES (100001, 'user_001', 12345, 1, 5000, TO_TIMESTAMP('2023-02-15 16:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100002, 'user_002', 12346, 2, 4000, TO_TIMESTAMP('2023-02-15 15:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100003, 'user_003', 12347, 3, 3000, TO_TIMESTAMP('2023-02-15 14:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100004, 'user_001', 12347, 4, 2000, TO_TIMESTAMP('2023-02-15 13:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100005, 'user_002', 12348, 5, 1000, TO_TIMESTAMP('2023-02-15 12:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100006, 'user_001', 12348, 1, 1000, TO_TIMESTAMP('2023-02-15 11:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100007, 'user_003', 12347, 4, 2000, TO_TIMESTAMP('2023-02-15 10:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1); INSERT INTO orders_pay VALUES (2001, 100001, 1, TO_TIMESTAMP('2023-02-15 17:40:56')), (2002, 100002, 1, TO_TIMESTAMP('2023-02-15 17:40:56')), (2003, 100003, 0, TO_TIMESTAMP('2023-02-15 17:40:56')), (2004, 100004, 0, TO_TIMESTAMP('2023-02-15 17:40:56')), (2005, 100005, 0, TO_TIMESTAMP('2023-02-15 18:40:56')), (2006, 100006, 0, TO_TIMESTAMP('2023-02-15 18:40:56')), (2007, 100007, 0, TO_TIMESTAMP('2023-02-15 18:40:56')); INSERT INTO product_catalog VALUES (1, 'phone_aaa'), (2, 'phone_bbb'), (3, 'phone_ccc'), (4, 'phone_ddd'), (5, 'phone_eee');NotaAs tabelas criadas neste tutorial são tabelas Paimon do tipo append-only, sem chaves primárias. Elas oferecem melhor desempenho de gravação em lote do que tabelas com chave primária, mas não suportam operações de atualização baseadas em chave.
O resultado da execução conterá várias subabas. A mensagem
The following statement has been executed successfully!confirma que a instrução DDL correspondente foi executada com êxito.Instruções DML, como INSERT, retornam um JobId, indicando que há um job Flink em execução na sessão. Clique em Flink UI à esquerda da aba Results para verificar o status dessas instruções SQL. Aguarde alguns segundos até a conclusão.
-
Explore os dados da tabela ODS.
No editor da aba Scripts, insira as seguintes instruções SQL e clique em Run à esquerda.
SELECT count(*) as order_count FROM `my_catalog`.`order_dw`.`orders`; SELECT count(*) as pay_count FROM `my_catalog`.`order_dw`.`orders_pay`; SELECT * FROM `my_catalog`.`order_dw`.`product_catalog`;Essas instruções SQL também rodam na sessão Flink. Visualize os resultados nas respectivas páginas de resultados.
O valor de
order_countpara a tabela orders é 7. O valor depay_countpara a tabela orders_pay também é 7. Já a tabela product_catalog contém 5 registros comproduct_idde 1 a 5 e valores decatalog_nameiguais a phone_aaa, phone_bbb, phone_ccc, phone_ddd e phone_eee, respectivamente.
Etapa 2: Criar tabelas DWD e DWS
No editor da aba Scripts, insira as seguintes instruções SQL e clique em Run à esquerda.
USE `my_catalog`.`order_dw`;
CREATE TABLE dwd_orders (
order_id BIGINT,
order_user_id STRING,
order_shop_id BIGINT,
order_product_id BIGINT,
order_product_catalog_name STRING,
order_fee BIGINT,
order_create_time TIMESTAMP,
order_update_time TIMESTAMP,
order_state INT,
pay_id BIGINT,
pay_platform INT COMMENT 'platform 0: phone, 1: pc',
pay_create_time TIMESTAMP
) WITH (
'sink.parallelism' = '2'
);
CREATE TABLE dws_users (
user_id STRING,
ds STRING,
total_fee BIGINT COMMENT 'Total amount paid on the current day'
) WITH (
'sink.parallelism' = '2'
);
CREATE TABLE dws_shops (
shop_id BIGINT,
ds STRING,
total_fee BIGINT COMMENT 'Total amount paid on the current day'
) WITH (
'sink.parallelism' = '2'
);
As tabelas criadas aqui também são do tipo append-only do Paimon. Quando uma tabela Paimon funciona como sink do Flink, o paralelismo não é inferido automaticamente. Defina explicitamente o paralelismo para evitar possíveis erros.
Etapa 3: Criar e implantar jobs DWD e DWS
-
Crie e implante a implantação DWD.
-
Crie a implantação de atualização da tabela DWD.
Na página , crie um blank batch draft chamado
dwd_orders. Recomendamos selecionar uma versão de engine com a tag recomendada. Use o dialeto Flink SQL padrão. Copie a seguinte instrução SQL para o editor. Como a tabela DWD é do tipo append-only do Paimon, utilize uma instruçãoINSERT OVERWRITEpara sobrescrever toda a tabela.INSERT OVERWRITE my_catalog.order_dw.dwd_orders SELECT o.order_id, o.user_id, o.shop_id, o.product_id, c.catalog_name, o.buy_fee, o.create_time, o.update_time, o.state, p.pay_id, p.pay_platform, p.create_time FROM my_catalog.order_dw.orders as o, my_catalog.order_dw.product_catalog as c, my_catalog.order_dw.orders_pay as p WHERE o.product_id = c.product_id AND o.order_id = p.order_id No canto superior direito da página, clique em Deploy e, em seguida, clique em OK para implantar o job
dwd_orders.
-
-
Crie e implante as implantações DWS.
-
Crie as implantações de atualização das tabelas DWS.
Seguindo os mesmos passos descritos em Criar a implantação de atualização da tabela DWD, crie dois rascunhos em lote chamados
dws_shopsedws_users. Copie as seguintes instruções SQL para seus respectivos editores.INSERT OVERWRITE my_catalog.order_dw.dws_shops SELECT order_shop_id, DATE_FORMAT(pay_create_time, 'yyyyMMdd') as ds, SUM(order_fee) as total_fee FROM my_catalog.order_dw.dwd_orders WHERE pay_id IS NOT NULL AND order_fee IS NOT NULL GROUP BY order_shop_id, DATE_FORMAT(pay_create_time, 'yyyyMMdd');INSERT OVERWRITE my_catalog.order_dw.dws_users SELECT order_user_id, DATE_FORMAT(pay_create_time, 'yyyyMMdd') as ds, SUM(order_fee) as total_fee FROM my_catalog.order_dw.dwd_orders WHERE pay_id IS NOT NULL AND order_fee IS NOT NULL GROUP BY order_user_id, DATE_FORMAT(pay_create_time, 'yyyyMMdd'); No canto superior direito da página, clique em Deploy e, em seguida, clique em OK para implantar os jobs
dws_shopsedws_users.
-
Etapa 4: Iniciar e visualizar implantações DWD e DWS
-
Inicie a implantação DWD e visualize seus dados.
-
Acesse a página . Selecione BATCH na lista suspensa. Localize a implantação
dwd_orderse clique em Start na coluna Actions.Uma instância de job com estado STARTING aparecerá na lista de instâncias de job em lote.
Quando o status da instância mudar para FINISHED, o processamento de dados estará concluído.
-
Explore os dados resultantes.
No editor da aba Scripts, insira a seguinte instrução SQL e clique em Run à esquerda para consultar os dados na tabela DWD.
SELECT * FROM `my_catalog`.`order_dw`.`dwd_orders`;O resultado é mostrado na figura a seguir.

-
-
Inicie as implantações DWS e visualize seus dados.
Na página , selecione BATCH na lista suspensa. Localize as implantações
dws_shopsedws_userse clique em Start na coluna Actions de cada uma.-
No editor da aba Scripts, insira as seguintes instruções SQL e clique em Run à esquerda para consultar os dados nas tabelas DWS.
SELECT * FROM `my_catalog`.`order_dw`.`dws_shops`; SELECT * FROM `my_catalog`.`order_dw`.`dws_users`;Os resultados são mostrados nas figuras a seguir.

Etapa 5: Construir um pipeline em lote com orquestração de tarefas
Nesta etapa, você orquestrará as implantações em um fluxo de trabalho, permitindo que sejam acionadas juntas e executadas em sequência.
-
Crie um fluxo de trabalho.
No painel de navegação à esquerda, escolha e clique em Create Workflow.
No painel exibido, insira
wf_orderscomo nome. Mantenha o tipo de agendamento como padrão (Manual Scheduling), selecionedefault-queuepara a Resource Queue e clique em Create para abrir o editor de fluxo de trabalho.-
Edite o fluxo de trabalho.
Clique no nó inicial, nomeie-o como
v_dwd_orderse selecione a implantaçãodwd_orderspara ele.Clique em Add Task para criar um nó chamado
v_dws_shops. Selecione a implantaçãodws_shopse definav_dwd_orderscomo nó upstream.Clique novamente em Add Task para criar um nó chamado
v_dws_users. Selecione a implantaçãodws_userse definav_dwd_orderscomo nó upstream.No canto superior direito, clique em Save e depois em OK.
-
Acione o fluxo de trabalho manualmente.
NotaTambém é possível configurar o fluxo para rodar periodicamente. Na página Workflows, clique em Edit Workflow à direita do fluxo desejado e altere o modo de agendamento para Periodic Scheduling. Para mais detalhes, consulte Orquestração de tarefas (Public Preview).
-
Antes de acionar o fluxo, insira novos dados nas tabelas ODS para validar os resultados da execução.
No editor da aba Scripts, insira as seguintes instruções SQL e clique em Run à esquerda.
USE `my_catalog`.`order_dw`; INSERT INTO orders VALUES (100008, 'user_001', 12346, 1, 10000, TO_TIMESTAMP('2023-02-15 17:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100009, 'user_002', 12347, 2, 20000, TO_TIMESTAMP('2023-02-15 18:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1), (100010, 'user_003', 12348, 3, 30000, TO_TIMESTAMP('2023-02-15 19:40:56'), TO_TIMESTAMP('2023-02-15 18:42:56'), 1); INSERT INTO orders_pay VALUES (2008, 100008, 1, TO_TIMESTAMP('2023-02-15 20:40:56')), (2009, 100009, 1, TO_TIMESTAMP('2023-02-15 20:40:56')), (2010, 100010, 1, TO_TIMESTAMP('2023-02-15 20:40:56'));Clique em Flink UI à esquerda da aba Results para monitorar o status do job.
-
Na página , localize o fluxo criado na etapa anterior, clique em Execute na coluna Actions e confirme clicando em OK para iniciar o fluxo.
Clique no nome do fluxo para acessar a página Workflow Instance List and Details, onde é possível ver a lista de instâncias do fluxo de trabalho.
Clique no ID da instância em execução para abrir a página de detalhes e acompanhar o status de cada nó. Aguarde a conclusão de todo o fluxo.
-
-
Visualize os resultados da execução do fluxo.
-
No editor da aba Scripts, insira as seguintes instruções SQL e clique em Run à esquerda.
SELECT * FROM `my_catalog`.`order_dw`.`dws_shops`; SELECT * FROM `my_catalog`.`order_dw`.`dws_users`; -
Verifique os resultados finais.
Observe que o fluxo processou os novos dados da camada ODS e os gravou nas tabelas DWS.

-
Para entender melhor os princípios e o ajuste de configuração do processamento em lote do Flink, consulte o Guia de ajuste do processamento em lote do Flink.
Para construir um data lakehouse de streaming usando Flink e Paimon, consulte Construir um data lakehouse de streaming com Paimon e StarRocks.
Além de desenvolver jobs Flink no console do Realtime Compute for Apache Flink, você também pode desenvolvê-los localmente. Para mais informações, consulte Desenvolver com a extensão do VS Code.