Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Processamento em lote do Flink

Última atualização: Jul 03, 2026

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.

image

Preparações

  1. Crie um script.

    Na aba Scripts, crie um catálogo contendo bancos de dados e tabelas e insira dados de amostra nessas tabelas.

  2. Crie um catálogo Paimon.

    1. 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 warehouse estiver 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 warehouse estiver 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

    2. 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

Nota

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.

  1. 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');
    Nota

    As 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.

  2. 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_count para a tabela orders é 7. O valor de pay_count para a tabela orders_pay também é 7. Já a tabela product_catalog contém 5 registros com product_id de 1 a 5 e valores de catalog_name iguais 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'
);
Nota

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

  1. Crie e implante a implantação DWD.

    1. Crie a implantação de atualização da tabela DWD.

      Na página Development > ETL, 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ção INSERT OVERWRITE para 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
    2. No canto superior direito da página, clique em Deploy e, em seguida, clique em OK para implantar o job dwd_orders.

  2. Crie e implante as implantações DWS.

    1. 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_shops e dws_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');
    2. No canto superior direito da página, clique em Deploy e, em seguida, clique em OK para implantar os jobs dws_shops e dws_users.

Etapa 4: Iniciar e visualizar implantações DWD e DWS

  • Inicie a implantação DWD e visualize seus dados.

    1. Acesse a página O&M > Deployments. Selecione BATCH na lista suspensa. Localize a implantação dwd_orders e 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.

    2. 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.

      image

  • Inicie as implantações DWS e visualize seus dados.

    1. Na página O&M > Deployments, selecione BATCH na lista suspensa. Localize as implantações dws_shops e dws_users e clique em Start na coluna Actions de cada uma.

    2. 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.

      image.png image.png

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.

  1. Crie um fluxo de trabalho.

    1. No painel de navegação à esquerda, escolha O&M > Workflows e clique em Create Workflow.

    2. No painel exibido, insira wf_orders como nome. Mantenha o tipo de agendamento como padrão (Manual Scheduling), selecione default-queue para a Resource Queue e clique em Create para abrir o editor de fluxo de trabalho.

    3. Edite o fluxo de trabalho.

      1. Clique no nó inicial, nomeie-o como v_dwd_orders e selecione a implantação dwd_orders para ele.

      2. Clique em Add Task para criar um nó chamado v_dws_shops. Selecione a implantação dws_shops e defina v_dwd_orders como nó upstream.

      3. Clique novamente em Add Task para criar um nó chamado v_dws_users. Selecione a implantação dws_users e defina v_dwd_orders como nó upstream.

      4. No canto superior direito, clique em Save e depois em OK.

  2. Acione o fluxo de trabalho manualmente.

    Nota

    També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).

    1. 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.

    2. Na página O&M > Workflows, 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.

  3. Visualize os resultados da execução do fluxo.

    1. 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`;
    2. Verifique os resultados finais.

      Observe que o fluxo processou os novos dados da camada ODS e os gravou nas tabelas DWS.

      image.png image.png

  4. Documentos relacionados