Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Crie um data warehouse em tempo real com o Hologres

Última atualização: Aug 13, 2026

Este guia demonstra como construir um data warehouse em tempo real usando o Realtime Compute for Apache Flink e o Hologres. A solução combina o processamento de fluxo do Flink com recursos exclusivos do Hologres — como binary logging, armazenamento híbrido linha-coluna e isolamento rigoroso de recursos — para lidar com volumes crescentes de dados e atender às demandas de negócios em tempo real.

Contexto

Com a digitalização dos negócios, a demanda por dados mais recentes cresce rapidamente. Os data warehouses offline tradicionais, projetados para processamento em lote de grandes volumes de dados, já não são suficientes. Muitos cenários modernos exigem processamento, armazenamento e análise de dados em tempo real. Embora a metodologia para construção de data warehouses offline com arquiteturas em camadas, como ODS, DWD e DWS, seja bem estabelecida, faltava uma estrutura clara para seus equivalentes em tempo real. O uso de um data warehouse em tempo real permite um fluxo de dados eficiente e imediato entre cada camada.

Caso de uso

Este guia usa uma plataforma de e-commerce como exemplo para demonstrar a construção de um data warehouse em tempo real integrando Flink e Hologres. Essa abordagem permite processar e limpar dados em tempo real, fornecer dados em camadas reutilizáveis para aplicações downstream e dar suporte a diversos cenários de negócios, incluindo painéis em tempo real (monitoramento de transações, análise comportamental, perfilamento de usuários) e recomendações personalizadas.

Arquitetura da solução

  1. Construa a camada ODS (operational data store): ingira dados do banco de dados de negócios em tempo real.

    O Flink sincroniza três tabelas de negócios do MySQL — orders (tabela de pedidos), orders_pay (tabela de pagamentos) e product_catalog (dicionário de categorias de product) — com o Hologres em tempo real. Essas tabelas formam a camada ODS.

  2. Construa a camada DWD (data warehouse detail): crie uma tabela wide em tempo real.

    O Flink une as tabelas ODS em tempo real para criar uma tabela wide na camada DWD.

  3. Construa a camada DWS (data warehouse service): calcule métricas em tempo real.

    O Flink consome alterações de binary logging da tabela wide orientado a eventos para agregar métricas em tabelas específicas de usuários e lojas na camada DWS.

  4. Atenda a consultas de aplicações com o Hologres.

    • Consulte tabelas de métricas agregadas na camada DWS, com suporte a milhões de requisições por segundo (RPS).

    • Execute consultas OLAP na tabela wide DWD ou exiba relatórios em tempo real baseados nesses dados, com respostas em segundos.

Benefícios e capacidades principais

Esta solução oferece os seguintes benefícios:

  • Atualizações eficientes e consultas imediatas: o Hologres suporta atualizações, correções e acesso imediato a consultas em cada camada de dados. Isso resolve desafios comuns em data warehouses em tempo real tradicionais, nos quais dados intermediários são difíceis de consultar, atualizar e corrigir.

  • Camadas de dados e reutilização: cada camada de dados no Hologres pode atender independentemente a aplicações externas. Isso possibilita a reutilização eficiente de dados, criando um data warehouse em camadas reutilizável.

  • Arquitetura simplificada e maior eficiência: usar o Flink SQL para construir o pipeline ETL em tempo real e armazenar todas as camadas de dados (ODS, DWD e DWS) no Hologres simplifica a arquitetura e melhora a eficiência do processamento.

A solução baseia-se em três capacidades principais do Hologres, conforme mostra a tabela a seguir.

Capacidade principal

Descrição

Binary logging

O Hologres fornece binary logging, permitindo que o Flink leia alterações de dados em tempo real. Assim, o Hologres atua como source de streaming para jobs do Flink.

Hybrid row-column storage

O Hologres suporta um formato de armazenamento híbrido no qual uma única tabela armazena dados nos formatos orientado a linhas e orientado a colunas com forte consistência. Isso permite que uma tabela intermediária sirva como source do Flink, tabela de dimensão para point queries e temporal joins, além de fonte de dados para outras aplicações, como consultas OLAP ou service online.

Isolamento rigoroso de recursos

Cargas elevadas em uma instância do Hologres podem afetar o desempenho de point queries em camadas de dados intermediárias. O Hologres oferece isolamento rigoroso de recursos por meio de read/write splitting for primary and secondary instances (shared storage) ou da virtual warehouse architecture. Isso garante que a ingestão de dados pelo Flink a partir de binary logs não afete os service online.

Observações de uso

  • Esta solução de data warehouse em tempo real é suportada apenas em instâncias dedicadas do Hologres.

  • Seu workspace do Realtime Compute for Apache Flink, sua instância do ApsaraDB RDS for MySQL e sua instância do Hologres devem estar na mesma VPC. Se estiverem em VPCs diferentes, conecte-as primeiro ou use endpoints públicos. Para mais informações, consulte How do I access other services across VPCs? e How do I access the Internet?.

  • Caso utilize um usuário RAM ou função RAM para acessar recursos do Realtime Compute for Apache Flink, Hologres e ApsaraDB RDS for MySQL, verifique se ele possui as permissões necessárias.

Etapa 1: Prepare o ambiente

Crie uma instância RDS for MySQL e prepare os dados

  1. Crie uma instância do ApsaraDB RDS for MySQL. Para mais informações, consulte Create an ApsaraDB RDS for MySQL instance.

    A instância do ApsaraDB RDS for MySQL deve estar na mesma VPC que seu workspace do Flink e a instância do Hologres.

  2. Crie um banco de dados e uma conta.

    Na instância de destino, crie um banco de dados chamado order_dw e uma conta padrão com permissões de leitura e escrita nesse banco. Para mais informações, consulte Create a database e Create an account.

  3. Prepare a source de dados MySQL CDC.

    1. Na página de detalhes da instância, clique em Log On to Database.

    2. Na página de login, insira o nome de usuário e a senha da conta do banco de dados criada e clique em Log On.

    3. Após fazer login, clique duas vezes no banco de dados order_dw para acessá-lo.

    4. No SQL Console, insira as instruções DDL a seguir para criar as tabelas de negócios e as instruções INSERT para preenchê-las.

      CREATE TABLE `orders` (
        order_id bigint not null primary key,
        user_id varchar(50) not null,
        shop_id bigint not null,
        product_id bigint not null,
        buy_fee numeric(20,2) not null,   
        create_time timestamp not null,
        update_time timestamp not null default now(),
        state int not null 
      );
      CREATE TABLE `orders_pay` (
        pay_id bigint not null primary key,
        order_id bigint not null,
        pay_platform int not null,
        create_time timestamp not null
      );
      CREATE TABLE `product_catalog` (
        product_id bigint not null primary key,
        catalog_name varchar(50) not null
      );
      -- Prepare data.
      INSERT INTO product_catalog VALUES(1, 'phone_aaa'),(2, 'phone_bbb'),(3, 'phone_ccc'),(4, 'phone_ddd'),(5, 'phone_eee');
      INSERT INTO orders VALUES
      (100001, 'user_001', 12345, 1, 5000.05, '2023-02-15 16:40:56', '2023-02-15 18:42:56', 1),
      (100002, 'user_002', 12346, 2, 4000.04, '2023-02-15 15:40:56', '2023-02-15 18:42:56', 1),
      (100003, 'user_003', 12347, 3, 3000.03, '2023-02-15 14:40:56', '2023-02-15 18:42:56', 1),
      (100004, 'user_001', 12347, 4, 2000.02, '2023-02-15 13:40:56', '2023-02-15 18:42:56', 1),
      (100005, 'user_002', 12348, 5, 1000.01, '2023-02-15 12:40:56', '2023-02-15 18:42:56', 1),
      (100006, 'user_001', 12348, 1, 1000.01, '2023-02-15 11:40:56', '2023-02-15 18:42:56', 1),
      (100007, 'user_003', 12347, 4, 2000.02, '2023-02-15 10:40:56', '2023-02-15 18:42:56', 1);
      INSERT INTO orders_pay VALUES
      (2001, 100001, 1, '2023-02-15 17:40:56'),
      (2002, 100002, 1, '2023-02-15 17:40:56'),
      (2003, 100003, 0, '2023-02-15 17:40:56'),
      (2004, 100004, 0, '2023-02-15 17:40:56'),
      (2005, 100005, 0, '2023-02-15 18:40:56'),
      (2006, 100006, 0, '2023-02-15 18:40:56'),
      (2007, 100007, 0, '2023-02-15 18:40:56');
  4. Clique em Execute e, em seguida, clique em Direct Execution.

Crie uma instância do Hologres e grupos de computação

  1. Adquira uma instância dedicada do Hologres. Para mais informações, consulte Purchase a Hologres instance.

    A instância do Hologres deve estar na mesma VPC que a instância do ApsaraDB RDS for MySQL. Para experimentar o isolamento rigoroso de recursos por meio da divisão de leitura/escrita, este exemplo usa o tipo de instância Virtual Warehouse e define Reserved Compute Resource como 64, permitindo a criação de grupos de computação adicionais.

  2. Após log on to the instance, crie um banco de dados e conceda permissões.

    Crie um banco de dados chamado order_dw (com o modelo de permissão simples ativado) e conceda privilégios administrativos ao usuário. Para detalhes sobre gerenciamento e autorização de bancos de dados, consulte Manage databases.

    Nota
    • Se não encontrar a conta na lista suspensa User, significa que ela ainda não foi adicionada à instância. Acesse a página User Management e adicione o usuário como SuperUser.

    • No Hologres V2.0 e versões posteriores, a extensão de binary logging vem habilitada por padrão. Não é necessário ativá-la manualmente.

  3. Crie um novo grupo de computação.

    Use grupos de computação diferentes para isolar recursos. Utilize o grupo de computação inicial init_warehouse para gravação de dados e o grupo read_warehouse_1 para atender consultas.

    Todos os recursos de computação reservados são alocados ao grupo de computação inicial init_warehouse por padrão. Reduza seus recursos antes de criar um novo grupo de computação. Para mais informações, consulte Create a new compute group instance.

    1. Acesse Security > Compute Group Management e confirme o nome da instância.

    2. Na linha do grupo de computação init_warehouse, clique em Modify Configuration na coluna Actions. Reduza os recursos alocados e clique em OK.

    3. Clique em Create Compute Group, crie um novo grupo de computação chamado read_warehouse_1 e clique em OK.

Crie um workspace do Flink e catálogos

  1. Crie um workspace do Flink. Para mais informações, consulte Activate Realtime Compute for Apache Flink.

    O workspace do Flink deve estar na mesma VPC que as instâncias do ApsaraDB RDS for MySQL e do Hologres.

  2. Faça login no console do Realtime Compute for Apache Flink e clique em console na coluna Actions do seu workspace.

  3. Crie um cluster de sessão para fornecer um ambiente de execução para criação de catálogos e execução de scripts. Para mais informações, consulte Step 1: Create a session cluster.

  4. Crie um catálogo do Hologres.

    Na página Development > Scripts, na aba Scripts, copie o código a seguir, substitua os valores de placeholder, selecione o código e clique em Run. Isso utiliza o cluster de sessão criado como ambiente de execução.

    CREATE CATALOG dw WITH (
      'type' = 'hologres',
      'endpoint' = '< ENDPOINT>', 
      'username' = 'BASIC$flinktest',
      'password' = '${secret_values. holosecrect}',
      'dbname' = 'order_dw@init_warehouse', -- Specify the database name and connect to the init_warehouse compute group.
      'binlog' = 'true', -- You can set default WITH options for source, dimension, and result tables when creating the catalog. Tables created under this catalog inherit these defaults.
      'sdkMode' = 'jdbc', -- The jdbc mode is recommended.
      'cdcmode' = 'true',
      'connectionpoolname' = 'the_conn_pool',
      'ignoredelete' = 'true',  -- Required for wide-table merge to prevent retractions.
      'partial-insert. enabled' = 'true', -- Required for wide-table merge to enable partial column updates.
      'mutateType' = 'insertOrUpdate', -- Required for wide-table merge to enable partial column updates.
      'table_property. binlog. level' = 'replica', -- You can also pass persistent Hologres table properties when creating the catalog. Tables created later will have binary logging enabled by default.
      'table_property. binlog. ttl' = '259200'
    );

    Modifique os parâmetros a seguir com as informações reais do seu service Hologres.

    Parâmetro

    Descrição

    Observações

    endpoint

    O endpoint da sua instância do Hologres.

    Na página de detalhes da instância do Hologres, obtenha o nome de domínio da VPC especificada. Para mais informações sobre nomes de domínio, consulte Endpoints.

    username

    Escolha uma das opções:

    • O nome de usuário de uma conta personalizada deve ter o formato BASIC$< user_name>.

    • O AccessKey ID da sua conta Alibaba Cloud ou usuário RAM.

    • O usuário configurado deve ter acesso ao banco de dados correspondente do Hologres. Para detalhes, consulte Hologres permission model e Manage users.

    • Este exemplo usa uma conta personalizada chamada BASIC$flinktest e define sua senha usando uma variável de projeto chamada holosecrect para evitar riscos de segurança ao armazenar senhas em texto simples. Para mais informações, consulte Project variables.

    password

    • A senha da conta personalizada.

    • O AccessKey secret da sua conta Alibaba Cloud ou usuário RAM.

    Nota

    Ao criar um catálogo, defina opções WITH padrão para tabelas de source, dimensão e resultado. Também é possível definir propriedades padrão para tabelas físicas do Hologres, como os parâmetros que começam com table_property. Para mais informações, consulte Manage Hologres catalogs e Hologres connector for real-time data warehouses (WITH parameters).

  5. Crie um catálogo MySQL.

    Copie o código a seguir na aba Scripts, modifique os valores dos parâmetros, selecione o código e clique em Run. Isso utiliza o cluster de sessão criado como ambiente de execução.

    CREATE CATALOG mysqlcatalog WITH(
      'type' = 'mysql',
      'hostname' = '< hostname>',
      'port' = '< port>',
      'username' = '< username>',
      'password' = '${secret_values. mysql_pw}',
      'default-database' = 'order_dw'
    );

    Modifique os parâmetros a seguir com as informações reais do seu service MySQL.

    Parâmetro

    Descrição

    hostname

    O endereço IP ou nome de host do seu banco de dados MySQL. Na página de informações básicas do banco de dados, clique em View Connection Details na área Network Type para obter o endpoint interno.

    port

    O número da porta do seu service de banco de dados MySQL. O valor padrão é 3306.

    username

    O nome de usuário do seu service de banco de dados MySQL.

    password

    A senha do seu service de banco de dados MySQL.

    Este exemplo usa uma variável chamada mysql_pw para a senha, evitando exposição em texto simples. Para mais informações, consulte Manage variables.

Etapa 2: Construa o data warehouse em tempo real

Construa a camada ODS: Ingira dados de negócios

Com a instrução CREATE DATABASE AS (CDAS) baseada em catálogo, crie a camada ODS em uma única etapa. A camada ODS geralmente serve como source de eventos para jobs de streaming, e não diretamente para consultas OLAP ou point queries. Habilitar o binary logging é suficiente para essa finalidade. O binary logging é uma capacidade central do Hologres. O conector do Hologres também suporta um modo completo mais incremental: ele lê um snapshot completo primeiro e depois consome binary logs incrementalmente.

  1. Crie o job de sincronização CDAS da ODS.

    1. Na página Development > ETL, crie um novo rascunho de fluxo SQL chamado ODS e copie o código a seguir no editor SQL.

      -- The table_property. binlog. level parameter was set when creating the catalog, so all tables created by CDAS have binary logging enabled.
      CREATE DATABASE IF NOT EXISTS dw. order_dw   
      AS DATABASE mysqlcatalog. order_dw INCLUDING all tables -- You can select the upstream tables to ingest as needed.
      /*+ OPTIONS('server-id'='8001-8004') */;   -- Specify the server-id range for the mysql-cdc instance.
      Nota
      • Por padrão, este exemplo sincroniza dados para o schema Public do banco de dados order_dw. Também é possível sincronizar dados para um schema específico no banco de dados Hologres de destino. Para mais informações, consulte Use a Hologres catalog as the destination in a CREATE DATABASE AS... statement. Após especificar um schema, o formato do nome da tabela muda ao usar o catálogo. Para detalhes, consulte Use a Hologres catalog.

      • Se o schema de uma tabela de origem mudar, o schema da tabela de resultado só será atualizado quando ocorrer uma alteração de dados (exclusão, inserção ou atualização) na tabela de origem.

    2. No canto superior direito, clique em Deploy para implantar o job.

    3. No painel de navegação à esquerda, escolha O&M > Deployments. Na linha do job ODS recém-implantado, clique em Start na coluna Actions. Selecione Initial Mode e clique em Start.

  2. Carregue dados no grupo de computação.

    Um table group é um portador de dados no Hologres. Ao usar o grupo de computação read_warehouse_1 para consultar dados de um table group no banco de dados order_dw, como order_dw_tg_default (para criar um table group, consulte Table Group Management), o table group order_dw_tg_default é carregado para o grupo de computação read_warehouse_1. Isso permite usar o grupo de computação init_warehouse para gravar dados e o grupo read_warehouse_1 para consultas de service.

    Na página de desenvolvimento do HoloWeb, clique em SQL Editor. Confirme o nome da instância e do banco de dados e execute os comandos a seguir. Para mais informações, consulte Create a new compute group instance. Após o carregamento, verifique que read_warehouse_1 carregou os dados do table group order_dw_tg_default.

    -- List table groups in the current database.
    SELECT tablegroup_name FROM hologres. hg_table_group_properties GROUP BY tablegroup_name;
    -- Load a table group into a compute group.
    CALL hg_table_group_load_to_warehouse ('order_dw. order_dw_tg_default', 'read_warehouse_1', 1);
    -- Check the table groups loaded into the compute group.
    select * from hologres. hg_warehouse_table_groups;
  3. No canto superior direito, alterne o grupo de computação para read_warehouse_1. Consultas e análises subsequentes usarão este grupo de computação.

    No canto superior direito da página do HoloWeb, selecione read_warehouse_1 na lista suspensa de grupos de computação.

  4. Na página SQL Editor, execute os comandos a seguir para visualizar os dados sincronizados do MySQL para as três tabelas do Hologres.

    -- Query data from the orders table.
    SELECT * FROM orders;
    -- Query data from the orders_pay table.
    SELECT * FROM orders_pay;
    -- Query data from the product_catalog table.
    SELECT * FROM product_catalog;

    O resultado da consulta da tabela product_catalog contém duas colunas, product_id (1 a 5) e catalog_name (phone_aaa, phone_bbb, phone_ccc, phone_ddd, phone_eee), totalizando 5 registros. Isso indica que os dados foram sincronizados com sucesso para o Hologres.

Construa a camada DWD: Crie uma tabela wide em tempo real

Esta etapa utiliza a capacidade de atualização parcial de colunas do conector do Hologres. Expresse atualizações parciais com INSERT DML. O job consulta múltiplas tabelas de dimensão usando point queries de alto desempenho, habilitadas pelo armazenamento em linhas e pelo armazenamento híbrido linha-coluna do Hologres. Com o isolamento rigoroso de recursos, as cargas de trabalho de escrita, leitura e análise não interferem umas nas outras.

  1. Use o recurso de catálogo do Flink para criar a tabela wide da camada DWD dwd_orders no Hologres.

    Na página Development > Scripts, copie o código a seguir na aba Scripts, selecione o código e clique em Run.

    -- Wide table columns must be nullable because different streams write to the same result table, and any column can be null.
    CREATE TABLE dw. order_dw. dwd_orders (
      order_id bigint not null,
      order_user_id string,
      order_shop_id bigint,
      order_product_id bigint,
      order_product_catalog_name string,
      order_fee numeric(20,2),
      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,
      PRIMARY KEY(order_id) NOT ENFORCED
    );
    -- You can modify Hologres physical table properties through the catalog.
    ALTER TABLE dw. order_dw. dwd_orders SET (
      'table_property. binlog. ttl' = '604800' -- Change the binary log TTL to one week.
    );
  2. Consuma alterações de binary logging das tabelas da camada ODS orders e orders_pay em tempo real.

    Na página Development > ETL, crie um novo rascunho de fluxo SQL chamado DWD. Copie o código a seguir no editor SQL e, em seguida, faça o Deploy e Start do job. Este job SQL une as tabelas orders e product_catalog usando um temporal join e grava o resultado na tabela dwd_orders, enriquecendo os dados em tempo real.

    BEGIN STATEMENT SET;
    INSERT INTO dw. order_dw. dwd_orders 
     (
       order_id,
       order_user_id,
       order_shop_id,
       order_product_id,
       order_fee,
       order_create_time,
       order_update_time,
       order_state,
       order_product_catalog_name
     ) SELECT o.*, dim. catalog_name 
       FROM dw. order_dw. orders as o
       LEFT JOIN dw. order_dw. product_catalog FOR SYSTEM_TIME AS OF proctime() AS dim
       ON o. product_id = dim. product_id;
    INSERT INTO dw. order_dw. dwd_orders 
      (pay_id, order_id, pay_platform, pay_create_time)
       SELECT * FROM dw. order_dw. orders_pay;
    END;
  3. Visualize os dados na tabela wide dwd_orders.

    Conecte-se à instância do Hologres na página de desenvolvimento do HoloWeb, faça login no banco de dados de destino e execute o comando a seguir no SQL Editor.

    SELECT * FROM dwd_orders;

    Após a execução bem-sucedida, o resultado da consulta retorna dados da tabela wide dwd_orders, incluindo campos como order_id, order_user_id, order_shop_id, order_product_id, order_product_catalog_name, order_fee, order_create_time e order_update_time.

Construa a camada DWS: Calcule métricas em tempo real

  1. Use o recurso de catálogo do Flink para criar as tabelas agregadas da camada DWS dws_users e dws_shops no Hologres.

    Na página Development > Scripts, copie o código a seguir na aba Scripts, selecione o código e clique em Run.

    -- User-dimension aggregate table.
    CREATE TABLE dw. order_dw. dws_users (
      user_id string not null,
      ds string not null,
      paied_buy_fee_sum numeric(20,2) not null comment 'Total amount of payments completed on the day',
      primary key(user_id, ds) NOT ENFORCED
    );
    -- Shop-dimension aggregate table.
    CREATE TABLE dw. order_dw. dws_shops (
      shop_id bigint not null,
      ds string not null,
      paied_buy_fee_sum numeric(20,2) not null comment 'Total amount of payments completed on the day',
      primary key(shop_id, ds) NOT ENFORCED
    );
  2. Consuma a tabela wide da camada DWD dw. order_dw. dwd_orders em tempo real, realize agregações no Flink e grave os resultados finais nas tabelas DWS no Hologres.

    Na página Development > ETL, crie um novo rascunho de fluxo SQL chamado DWS. Copie o código a seguir no editor SQL e, em seguida, faça o Deploy e Start do job.

    BEGIN STATEMENT SET;
    INSERT INTO dw. order_dw. dws_users
      SELECT 
        order_user_id,
        DATE_FORMAT (pay_create_time, 'yyyyMMdd') as ds,
        SUM (order_fee)
        FROM dw. order_dw. dwd_orders c
        WHERE pay_id IS NOT NULL AND order_fee IS NOT NULL -- Both order and payment stream data have been written to the wide table.
        GROUP BY order_user_id, DATE_FORMAT (pay_create_time, 'yyyyMMdd');
    INSERT INTO dw. order_dw. dws_shops
      SELECT 
        order_shop_id,
        DATE_FORMAT (pay_create_time, 'yyyyMMdd') as ds,
        SUM (order_fee)
       FROM dw. order_dw. dwd_orders c
       WHERE pay_id IS NOT NULL AND order_fee IS NOT NULL -- Both order and payment stream data have been written to the wide table.
       GROUP BY order_shop_id, DATE_FORMAT (pay_create_time, 'yyyyMMdd');
    END;
  3. Visualize os resultados agregados na camada DWS. Os resultados são atualizados em tempo real conforme os dados upstream mudam.

    1. Visualize os dados no console do Hologres antes da alteração.

      Tabela dws_users

      SELECT * FROM dws_users;

      Após executar a consulta, o resultado retorna dados da tabela dws_users, que inclui três colunas: user_id, ds e paied_buy_fee_sum. No resultado de exemplo, a coluna user_id contém os valores user_001, user_002 e user_003; a coluna ds contém o valor 20230215; e a coluna paied_buy_fee_sum contém os valores 8000.08, 5000.05 e 5000.05, respectivamente. A coluna user_id identifica exclusivamente cada usuário.

      Tabela dws_shops

      SELECT * FROM dws_shops;

      O resultado da consulta mostra que a tabela dws_shops contém três colunas: shop_id (ID da loja), ds (partição de data) e paied_buy_fee_sum (valor do pagamento). Quatro linhas de dados de amostra são retornadas, confirmando que a tabela da camada DWS foi construída com sucesso.

    2. No console do RDS, insira um novo registro em cada uma das tabelas orders e orders_pay no banco de dados order_dw.

      INSERT INTO orders VALUES
      (100008, 'user_003', 12345, 5, 6000.02, '2023-02-15 09:40:56', '2023-02-15 18:42:56', 1);
      INSERT INTO orders_pay VALUES
      (2008, 100008, 1, '2023-02-15 19:40:56');
    3. Visualize os dados atualizados no console do Hologres.

      Tabela dwd_orders

      SELECT * FROM dwd_orders;

      Após executar a consulta, oito registros de pedidos são retornados da tabela dwd_orders, incluindo campos como order_id, order_user_id, order_shop_id, order_product_id, order_product_catalog_name, order_fee, order_create_time e order_update_time. O oitavo registro (order_id=100008, user_003, phone_eee, 6000.02) corresponde aos dados recém-gravados.

      Tabela dws_users

      SELECT * FROM dws_users;

      A consulta retorna três linhas com três colunas: user_id, ds e paied_buy_fee_sum: user_001 / 20230215 / 8000.08, user_002 / 20230215 / 5000.05 e user_003 / 20230215 / 11000.07. O maior valor total de pagamento é do user_003 (11000.07).

      Tabela dws_shops

      SELECT * FROM dws_shops;

      Após a execução da consulta, o resultado contém três colunas, shop_id, ds e paied_buy_fee_sum, com quatro linhas de dados mostrando os valores de taxa (11000.07, 4000.04, 7000.07 e 2000.02) para as lojas 12345, 12346, 12347 e 12348 em 20230215. As colunas shop_id e paied_buy_fee_sum são as principais colunas de métrica.

Analise os dados

Como o binary logging está habilitado, inspecione diretamente as alterações de dados. Se precisar realizar exploração ad-hoc de dados de negócios em resultados intermediários ou verificar a correção do cálculo final, cada camada desta solução é persistida, facilitando o exame do processo intermediário.

Análise em modo Streaming

Use o conector Print para confirmar se as mensagens enviadas para outras tabelas de resultado atendem às expectativas.

  1. Crie e inicie um job de análise de dados em streaming.

    Na página Development > ETL, crie um novo rascunho de fluxo SQL chamado Data-exploration. Copie o código a seguir no editor SQL e, em seguida, faça o Deploy e Start do job.

    -- Streaming mode profiling. Print output shows real-time data changes.
    CREATE TEMPORARY TABLE print_sink(
      order_id bigint not null,
      order_user_id string,
      order_shop_id bigint,
      order_product_id bigint,
      order_product_catalog_name string,
      order_fee numeric(20,2),
      order_create_time timestamp,
      order_update_time timestamp,
      order_state int,
      pay_id bigint,
      pay_platform int,
      pay_create_time timestamp,
      PRIMARY KEY(order_id) NOT ENFORCED
    ) WITH (
      'connector' = 'print'
    );
    INSERT INTO print_sink SELECT *
    FROM dw. order_dw. dwd_orders /*+ OPTIONS('startTime'='2023-02-15 12:00:00') */ -- Here, startTime is the generation time of the binary log.
    WHERE order_user_id = 'user_001';
  2. Visualize os resultados da análise de dados.

    Na página de detalhes de O&M > Deployments, clique no nome do job alvo. Na aba Logs, clique na aba Logs e, em seguida, clique em um link Path, ID em Running Task Managers. Na página Stdout, pesquise informações de log relacionadas a user_001.

    A saída do log exibe registros de alteração de dados CDC prefixados com +I (inserção), -U (antes da atualização) e +U (após a atualização), incluindo campos como order_id, order_user_id, order_shop_id, order_fee e order_create_time.

Análise em modo Batch

A análise em modo Batch não grava dados em uma tabela de resultado. Em vez disso, recupera o estado final dos dados no momento atual, permitindo visualizar os resultados diretamente na saída de depuração.

Na página Development > ETL, crie um rascunho de fluxo SQL, copie o código a seguir no editor SQL e clique em Debug. Para mais informações, consulte Debug a job.

O resultado da depuração na página de desenvolvimento de jobs do Flink é o seguinte.

SELECT *
FROM dw. order_dw. dwd_orders /*+ OPTIONS('binlog'='false') */ 
WHERE order_user_id = 'user_001' and order_create_time > '2023-02-15 12:00:00'; -- Batch mode supports filter pushdown to improve batch job execution efficiency.

Após a conclusão da depuração, o resultado da consulta retorna dois registros de pedidos que atendem às condições de filtro, incluindo campos como order_id, order_user_id, order_shop_id, order_product_id, order_product_catalog_name, order_fee, order_create_time e order_update_time. Isso verifica que o resultado da análise em modo Batch está conforme o esperado.

Etapa 3: Use o data warehouse em tempo real

A Etapa 2 mostrou como usar catálogos do Flink para construir um data warehouse em tempo real em camadas baseado em Flink e Hologres. As seções a seguir descrevem alguns cenários simples de aplicação.

Point query

Consulte as tabelas de métricas agregadas da camada DWS pela chave primária, com suporte a milhões de RPS.

Na página de desenvolvimento do HoloWeb, execute o SQL a seguir para consultar o valor de consumo de um usuário específico em uma data específica.

-- holo sql
SELECT * FROM dws_users WHERE user_id ='user_001' AND ds = '20230215';

O resultado da consulta retorna três colunas: user_id, ds e paied_buy_fee_sum (valor de consumo). O valor de consumo para user_001 em 20230215 é 8000.08.

Consulta OLAP

Execute consultas OLAP na tabela wide da camada DWD.

Na página de desenvolvimento do HoloWeb, execute o SQL a seguir para consultar os detalhes do pedido de um cliente específico em uma plataforma de pagamento específica em fevereiro de 2023.

-- holo sql
SELECT * FROM dwd_orders
WHERE order_create_time >= '2023-02-01 00:00:00'  and order_create_time < '2023-03-01 00:00:00'
AND order_user_id = 'user_001'
AND pay_platform = 0
ORDER BY order_create_time LIMIT 100;

O resultado da consulta retorna dois registros de pedidos com campos como order_id, order_user_id, order_shop_id, order_product_id, order_fee, order_create_time e order_update_time. Os IDs de pedido de exemplo são 100006 e 100004.

Relatórios em tempo real

Gere relatórios em tempo real com base nos dados da tabela wide da camada DWD. O armazenamento híbrido linha-coluna e as tabelas orientadas a colunas no Hologres fornecem excelentes capacidades de consulta OLAP, com respostas em segundos.

Na página de desenvolvimento do HoloWeb, execute o SQL a seguir para consultar o número total de pedidos e o valor total dos pedidos para cada categoria de product em fevereiro de 2023.

-- holo sql
SELECT
  TO_CHAR(order_create_time, 'YYYYMMDD') AS order_create_date,
  order_product_catalog_name,
  COUNT(*),
  SUM(order_fee)
FROM
  dwd_orders
WHERE
  order_create_time >= '2023-02-01 00:00:00'  and order_create_time < '2023-03-01 00:00:00'
GROUP BY
  order_create_date, order_product_catalog_name
ORDER BY
  order_create_date, order_product_catalog_name;

Após executar a instrução SQL, a aba Result exibe uma tabela com quatro colunas: order_create_date, order_product_catalog_name, count e sum. Por exemplo, os resultados para a data 20230215 mostram que a categoria phone_aaa teve 2 pedidos totalizando 6000.06, a categoria phone_bbb teve 1 pedido de 4000.04, e assim por diante para todas as cinco categorias.

Referências