Todos os produtos
Search
Central de documentação

Hologres:Crie um data warehouse em tempo real com Flink e Hologres

Última atualização: Jul 10, 2026

Combine o processamento em tempo real do Realtime Compute for Apache Flink com recursos do Hologres, como Binlog, armazenamento híbrido linha-coluna e forte isolamento de recursos, para construir um data warehouse em tempo real escalável. Essa arquitetura suporta volumes crescentes de dados e atende às demandas de negócios em tempo real.

Contexto

A crescente demanda por dados atualizados impulsiona as empresas além do processamento em lote offline tradicional, rumo ao processamento, armazenamento e análise de dados em tempo real. Embora o data warehousing offline siga uma metodologia bem definida com processamento em camadas (ODS > DWD > DWS > ADS) por meio de jobs agendados, uma estrutura comparável para data warehousing em tempo real ainda está emergindo. Esta solução aplica o conceito de Streaming Warehouse para criar um fluxo eficiente de dados em tempo real entre as camadas e resolver os desafios do nivelamento de dados em tempo real.

Cenário

Com uma plataforma de e-commerce como exemplo, este tópico mostra como construir um data warehouse em tempo real integrando o Flink ao Hologres. A configuração resultante permite o processamento e a limpeza de dados em tempo real, suporta consultas de aplicações upstream e estabelece camadas e reutilização de dados para cenários como painéis de transações, análise de comportamento, criação de perfis de usuários e recomendações personalizadas.

Arquitetura

image
  1. Construa a camada ODS: ingira dados de bancos de dados de negócios em tempo real.

    O MySQL contém três tabelas de negócios: orders (tabela de pedidos), orders_pay (tabela de pagamento de pedidos) e product_catalog (tabela de dicionário de categorias de produtos). O Flink sincroniza essas três tabelas com o Hologres em tempo real para formar a camada ODS.

  2. Construa a camada DWD: crie uma tabela ampla em tempo real.

    O Flink une as tabelas orders, product_catalog e orders_pay em tempo real para criar uma tabela ampla na camada DWD.

  3. Construa a camada DWS: calcule métricas em tempo real.

    O Flink consome o Binlog da tabela ampla em um processo orientado a eventos e agrega dados para criar tabelas de métricas para usuários e lojas na camada DWS.

  4. Atenda a consultas de aplicações por meio do Hologres.

    • As aplicações podem consultar as tabelas de métricas agregadas na camada DWS, suportando milhões de solicitações por segundo (RPS).

    • As aplicações podem executar análises OLAP na tabela ampla DWD ou exibir relatórios em tempo real com base nesses dados, com respostas em segundos.

Benefícios e capacidades principais

Principais benefícios:

  • Atualizações eficientes e consultas imediatas: o Hologres suporta atualizações eficientes, correções e consistência de leitura após gravação para dados em todas as camadas. Isso resolve uma limitação fundamental dos data warehouses em tempo real tradicionais, nos quais os dados das camadas intermediárias são difíceis de consultar, atualizar e corrigir.

  • Camadas e reutilização de dados: todas as camadas de dados no Hologres podem ser expostas a serviços externos de forma independente, permitindo camadas e reutilização eficientes de dados.

  • Arquitetura simplificada e maior eficiência: construir um pipeline ETL em tempo real com Flink SQL e armazenar dados das camadas ODS, DWD e DWS no Hologres simplifica a arquitetura e melhora a eficiência do processamento de dados.

Esta solução baseia-se em três capacidades principais do Hologres.

Capacidade principal

Descrição

Binlog

O Binlog do Hologres impulsiona o Flink a executar computações em tempo real e serve como source upstream para processamento de fluxo.

Armazenamento híbrido linha-coluna

O Hologres suporta um formato de armazenamento híbrido linha-coluna. Uma única tabela armazena dados em formatos orientados a linhas e a colunas com forte consistência. Tabelas intermediárias podem servir como tabelas source do Flink, como tabelas de dimensão para consultas pontuais e joins, e também serem consultadas por outras aplicações, como serviços OLAP e online.

Forte isolamento de recursos

Uma carga alta em uma instância do Hologres pode afetar o desempenho de consultas pontuais em camadas intermediárias. O Hologres oferece forte isolamento de recursos por meio de implantação com divisão de leitura/gravação para instâncias primárias e secundárias (armazenamento compartilhado) ou da arquitetura de instância de virtual warehouse. Isso garante que os jobs do Flink que extraem dados do Binlog do Hologres não afetem os serviços online.

Pré-requisitos

  • Apenas instâncias exclusivas do Hologres suportam esta solução de data warehouse em tempo real.

  • As instâncias do Realtime Compute for Apache Flink, RDS MySQL e Hologres devem estar na mesma VPC. Caso contrário, conecte as VPCs primeiro ou use endpoints públicos para acesso. Para mais informações, consulte Como acesso outros serviços entre VPCs? e Como acesso a Internet?.

  • Garanta que qualquer usuário RAM ou função RAM usada para acesso tenha as permissões necessárias para os recursos do Realtime Compute for Apache Flink, Hologres e RDS MySQL.

Etapa 1: Preparar recursos

Criar uma instância RDS MySQL e preparar uma fonte de dados

  1. Crie uma instância RDS MySQL. Para mais informações, consulte Criar uma instância ApsaraDB RDS for MySQL.

    A instância RDS MySQL deve estar na mesma VPC que o 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 gravação para esse banco de dados. Para detalhes, consulte Criar um banco de dados e Criar uma conta.

  3. Prepare a fonte de dados MySQL CDC.

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

    2. Na caixa de diálogo Connect to Instance, insira o nome de usuário e a senha da conta criada e clique em Sign in.

    3. Após fazer login, clique duas vezes no banco de dados order_dw na página da instância de banco de dados para alternar para ele.

    4. Na área SQL Console, escreva as instruções DDL para criar as três tabelas de negócios e as instruções para inserir dados.

      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 Upload e, em seguida, clique em Execute.

Criar uma instância do Hologres e uma virtual warehouse

  1. Crie uma instância exclusiva do Hologres. Para mais informações, consulte Comprar uma instância do Hologres.

    A instância do Hologres deve estar na mesma VPC que a instância RDS MySQL. Para experimentar a capacidade de forte isolamento de recursos do Hologres por meio da divisão de leitura/gravação, selecione Virtual Warehouse como tipo de instância e defina Reserved Computing Resources of Virtual Warehouse como 64. Isso permite criar uma nova virtual warehouse.

  2. Após fazer login na instância, crie um banco de dados e conceda permissões.

    Crie um banco de dados chamado order_dw (o Simple Permission Model deve estar ativado) e conceda permissões de administrador ao usuário. Para mais informações sobre como criar um banco de dados e conceder permissões, consulte DB Management.

    Nota
    • Se uma conta não aparecer na lista suspensa User Account, ela não foi adicionada à instância. Acesse a página Users para adicionar o usuário como SuperUser.

    • No Hologres V2.0 e versões posteriores, a extensão Binlog é ativada por padrão e não requer execução manual.

  3. Crie uma nova virtual warehouse.

    Use diferentes virtual warehouses para obter isolamento de recursos. Utilize a virtual warehouse inicial init_warehouse para gravar dados e a virtual warehouse read_warehouse_1 para atender consultas.

    Os recursos de computação reservados são totalmente alocados para a virtual warehouse inicial init_warehouse. Reduza seus recursos antes de criar uma nova. Para mais informações, consulte Criar uma nova instância de virtual warehouse.

    1. Clique em Security Center > Virtual Warehouse Management e confirme se o nome da instância está correto.

    2. Na linha da virtual warehouse existente init_warehouse, clique em Modify Configuration na coluna Actions. Reduza os recursos e clique em OK.

    3. Clique em Create Virtual Warehouse, crie uma nova virtual warehouse chamada read_warehouse_1 e clique em OK.

Criar um workspace do Flink e catálogos

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

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

  2. Faça login no console do Realtime Compute for Apache Flink e, na linha do workspace de destino, clique em Console na coluna Actions.

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

  4. Crie um catálogo do Hologres.

    Na aba Script da página Development > Scripts, copie o código a seguir para o editor de scripts. Modifique os valores dos parâmetros de destino, selecione o trecho desejado e clique em Run. No canto inferior direito, certifique-se de que o cluster de sessão criado esteja selecionado como ambiente de execução.

    CREATE CATALOG dw WITH (
      'type' = 'hologres',
      'endpoint' = '<ENDPOINT>', 
      'username' = 'BASIC$flinktest',
      'password' = '${secret_values.holosecret}',
      'dbname' = 'order_dw@init_warehouse', -- Specify the database name and connect to the init_warehouse virtual warehouse.
      'binlog' = 'true', -- When you create a catalog, you can set the WITH parameters for source, dimension, and result tables. These default parameters are automatically added when you use tables under this catalog.
      'sdkMode' = 'jdbc', -- The jdbc mode is recommended.
      'cdcmode' = 'true',
      'connectionpoolname' = 'the_conn_pool',
      'ignoredelete' = 'true',  -- This must be enabled for wide table merges to prevent retractions.
      'partial-insert.enabled' = 'true', -- This parameter must be enabled for wide table merges to allow partial column updates.
      'mutateType' = 'insertOrUpdate', -- This parameter must be enabled for wide table merges to allow partial column updates.
      'table_property.binlog.level' = 'replica', -- You can also pass persistent Hologres table properties when creating a catalog. Binlog is then enabled by default when you create tables.
      'table_property.binlog.ttl' = '259200'
    );

    Modifique os seguintes valores de parâmetro com as informações reais do seu service Hologres.

    Parâmetro

    Descrição

    Observações

    endpoint

    O endpoint da instância do Hologres.

    Obtenha o nome de domínio para o tipo de rede VPC especificado na página de detalhes da instância do Hologres. Para mais informações sobre nomes de domínio, consulte Endpoints.

    username

    Selecione uma das opções a seguir:

    • O nome de usuário da conta personalizada está no formato BASIC$<user_name>.

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

    • O usuário configurado aqui precisa de acesso ao banco de dados correspondente do Hologres. Para mais informações sobre permissões de banco de dados e gerenciamento de usuários do Hologres, consulte Modelo de permissão do Hologres e Gerenciamento de usuários.

    • 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 associados ao armazenamento em texto simples. Para mais informações, consulte Variáveis de projeto.

    password

    • A senha da conta personalizada.

    • O AccessKey Secret da sua conta Alibaba Cloud ou de um usuário RAM.

    Nota

    Ao criar um catálogo, defina parâmetros WITH padrão para tabelas source, de dimensão e de resultado. Também é possível definir propriedades padrão para a criação de tabelas físicas do Hologres, como os parâmetros que começam com table_property. Para mais informações, consulte Gerenciar catálogos do Hologres e Parâmetros WITH do conector Hologres.

  5. Crie um catálogo MySQL.

    Copie o código a seguir para o editor Script. Modifique os valores dos parâmetros de destino, selecione o trecho desejado e clique em Run no lado esquerdo da linha de código. No canto inferior direito, certifique-se de que o cluster de sessão criado esteja selecionado 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 seguintes valores de parâmetro com as informações reais do seu service MySQL.

    Parâmetro

    Descrição

    hostname

    O endereço IP ou nome de host do banco de dados MySQL. Obtenha o endereço interno clicando em View Details na área Network Type da página de informações básicas do banco de dados.

    port

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

    username

    O nome de usuário para o service de banco de dados MySQL.

    password

    A senha para o service de banco de dados MySQL.

    Este exemplo usa uma variável chamada mysql_pw para o valor da senha, evitando riscos como armazenamento em texto simples. Para mais informações, consulte Gerenciar variáveis.

Etapa 2: Construir o data warehouse em tempo real

Construir a camada ODS: Ingestão de dados em tempo real

Construa a camada ODS em uma única etapa usando a instrução CREATE DATABASE AS (CDAS) do catálogo. A camada ODS geralmente serve como driver de eventos para jobs de streaming, e não para consultas pontuais OLAP ou chave-valor; portanto, ativar o Binlog é suficiente. O conector do Hologres suporta um modo completo e incremental que lê primeiro todos os dados existentes e depois consome o Binlog incrementalmente.

  1. Crie um job de sincronização CDAS chamado ODS.

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

      CREATE DATABASE IF NOT EXISTS dw.order_dw   -- The table_property.binlog.level parameter was set when the catalog was created, so Binlog is enabled for all tables created via CDAS.
      AS DATABASE mysqlcatalog.order_dw INCLUDING all tables -- You can select which tables from the upstream database to ingest.
      /*+ OPTIONS('server-id'='8001-8004') */ ;   -- Specify the server-id range for the mysql-cdc instance.
      Nota
      • Por padrão, este exemplo sincroniza dados com o Public Schema do banco de dados order_dw. Também é possível sincronizar dados com um schema específico no banco de dados Hologres de destino. Para mais informações, consulte Usar como catálogo de destino para CDAS. Após especificar um schema, o formato do nome da tabela para uso do catálogo também mudará. Para mais informações, consulte Usar um catálogo do Hologres.

      • Alterações de schema em uma tabela source são propagadas para a tabela de resultado somente após uma operação DML subsequente (INSERT, UPDATE ou DELETE) ocorrer na tabela source.

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

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

  2. Carregue dados na virtual warehouse.

    Um table group é o portador de dados no Hologres. Ao usar a virtual warehouse read_warehouse_1 para consultar dados de um table group no banco de dados order_dw, como order_dw_tg_default (para etapas de criação, consulte Table Group Management), carregue order_dw_tg_default para read_warehouse_1. Isso permite usar a virtual warehouse init_warehouse para gravar dados e a virtual warehouse read_warehouse_1 para realizar 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 Criar uma nova instância de virtual warehouse. Após o carregamento, verifique que read_warehouse_1 carregou os dados do Table Group order_dw_tg_default.

    -- View the 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 virtual warehouse.
    CALL hg_table_group_load_to_warehouse ('order_dw.order_dw_tg_default', 'read_warehouse_1', 1);
    -- View the loading status of Table Groups for virtual warehouses.
    select * from hologres.hg_warehouse_table_groups;
  3. No canto superior direito, alterne a virtual warehouse para read_warehouse_1. Consultas e análises subsequentes usarão a virtual warehouse read_warehouse_1.

    A lista suspensa de virtual warehouse no canto superior direito mostra read_warehouse_1. O editor exibe a instrução de carregamento de table group executada CALL hg_table_group_load_to_warehouse ('order_dw.order_dw_tg_default', 'read_warehouse_1', 1); e a instrução de consulta select * from hologres.hg_warehouse_table_groups;.

  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 in the orders table.
    SELECT * FROM orders;
    --- Query data in the orders_pay table.
    SELECT * FROM orders_pay;
    --- Query data in the product_catalog table.
    SELECT * FROM product_catalog;

    Após executar a terceira consulta, a aba Result[3] mostra cinco linhas na tabela product_catalog, incluindo duas colunas: product_id (com valores de 1 a 5) e catalog_name (com valores phone_aaa, phone_bbb, phone_ccc, phone_ddd e phone_eee).

Construir a camada DWD: Criar uma tabela ampla em tempo real

A camada DWD aproveita a capacidade de atualização de coluna parcial do conector do Hologres, permitindo que instruções DML INSERT expressem semântica de atualização de coluna parcial. Esse processo depende de consultas pontuais de alto desempenho em tabelas de dimensão por meio do armazenamento de linhas e do armazenamento híbrido linha-coluna do Hologres, enquanto o forte isolamento de recursos garante que jobs de gravação, leitura e análise não interfiram entre si.

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

    Na aba Script da página Development > Scripts, copie o código a seguir no editor de scripts, selecione o trecho e clique em Run no lado esquerdo da linha de código.

    -- Wide table fields must be nullable because when different streams write to the same result table, any column can potentially have a null value.
    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 Binlog timeout to one week.
    );
  2. Consuma o Binlog das tabelas da camada ODS orders e orders_pay em tempo real.

    Na página Data Development>ETL, crie um job de streaming SQL chamado DWD, copie o código a seguir no editor SQL e, em seguida, Deploy e Start o job. Este job SQL une a tabela orders com a tabela de dimensão product_catalog, grava o resultado final na tabela dwd_orders e realiza enriquecimento de 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 ampla dwd_orders.

    Na página de desenvolvimento do HoloWeb, conecte-se à instância do Hologres e faça login no banco de dados de destino. Em seguida, execute o comando a seguir no editor SQL.

    SELECT * FROM dwd_orders;

    A tabela ampla dwd_orders inclui os seguintes campos: order_id, order_user_id, order_shop_id, order_product_id, order_product_catalog_name, order_fee, order_create_time, order_update_time, order_state, pay_id, pay_platform e pay_create_time. A consulta retorna 7 registros de pedidos.

Construir a camada DWS: Calcular métricas em tempo real

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

    Na aba Script da página Development > Scripts, copie o código a seguir no editor de scripts, selecione o trecho e clique em Run no lado esquerdo da linha de código.

    -- User-dimension aggregate metric 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 paid on the current day',
      primary key(user_id,ds) NOT ENFORCED
    );
    -- Shop-dimension aggregate metric 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 paid on the current day',
      primary key(shop_id,ds) NOT ENFORCED
    );
  2. Consuma dados da tabela ampla 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 Data Development>ETL, crie um novo job de streaming SQL chamado DWS, copie o código a seguir no editor SQL e, em seguida, Deploy e Start o 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 -- Data from both order and payment streams has 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 -- Data from both order and payment streams has been written to the wide table.
       GROUP BY order_shop_id, DATE_FORMAT (pay_create_time, 'yyyyMMdd');
    END;
  3. Visualize os resultados de agregação na camada DWS. Os resultados são atualizados em tempo real conforme os dados upstream mudam.

    1. No console do Hologres, visualize os dados antes da alteração.

      dws_users

      SELECT * FROM dws_users;

      O resultado da consulta contém três colunas: user_id (por exemplo, user_001, user_002, user_003), ds (por exemplo, 20230215) e paied_buy_fee_sum (por exemplo, 8000.08, 5000.05). A coluna user_id é o campo chave de associação.

      dws_shops

      SELECT * FROM dws_shops;

      A consulta retorna 4 registros com três colunas: shop_id, ds e paied_buy_fee_sum. Nos dados de exemplo, shop_id varia de 12345 a 12348, todos os valores de ds são 20230215 e os valores de paied_buy_fee_sum são 5000.05, 4000.04, 7000.07 e 2000.02, respectivamente.

    2. No console do RDS, insira um novo registro de dados em cada uma das tabelas orders e orders_pay do 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. No console do Hologres, visualize os dados após a alteração.

      dwd_orders

      SELECT * FROM dwd_orders;

      O resultado da execução mostra 8 registros totais de pedidos na tabela dwd_orders (order_id 100001–100008). Os campos incluem order_user_id, order_shop_id, order_product_id, order_product_catalog_name, order_fee, order_create_time e order_update_time. O registro recém-inserido para order_id=100008 tem um order_fee de 6000.02.

      dws_users

      SELECT * FROM dws_users;

      A consulta retorna 3 registros da tabela dws_users com colunas para user_id, ds e paied_buy_fee_sum. Os dados são: user_001 / 20230215 / 8000.08, user_002 / 20230215 / 5000.05 e user_003 / 20230215 / 11000.07. O valor agregado para user_003 é o mais alto, com 11000.07.

      dws_shops

      SELECT * FROM dws_shops;

      A consulta retorna quatro linhas com três colunas: shop_id, ds e paied_buy_fee_sum. Os valores de shop_id são 12345, 12346, 12347 e 12348; todos os valores de ds são 20230215; e os valores de paied_buy_fee_sum são 11000.07, 4000.04, 7000.07 e 2000.02, respectivamente.

Análise exploratória de dados

Com o Binlog ativado, inspecione diretamente as alterações nos dados. A persistência de dados em cada camada simplifica a análise exploratória ad hoc e a verificação de resultados.

Stream-mode profiling

Use o conector Print para ajudar a confirmar se as mensagens enviadas para outras tabelas de resultado estão conforme o esperado.

  1. Crie e inicie um job de análise exploratória em streaming.

    Na página Data Development>ETL, crie um job de streaming SQL chamado Data-exploration, copie o código a seguir no editor SQL e Deploy e Start o job.

    -- In stream-mode profiling, you can print to see 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 binlog.
    WHERE order_user_id = 'user_001';
  2. Visualize os resultados da análise exploratória de dados.

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

    NWoJuf*****]. secret: [CrxBZYHuTD*****], token: [CAISjgRxxx]
    end new OSSLogClient endTimeInMs:[1744628993550], costInMxxx
    [1744628993551], costInMs:[10 ms][OSSLogAppender:main] doSend cost time(ms):[59], current log queue size:[1], total received/discarded:[401/0],exceptionReceived/exceptionDiscarded:[0/0], total send:[400]
    [OSSLogAppender:main] doSend cost time(ms):[57], current log queue size:[2], total received/discarded:[502/0], exceptionReceived/exceptionDiscarded:[0/0], total send:[500]
    +I[100001, user_001, 12345, 1, phone_aaa, 5000.05, 2023-02-15T16:40:56, 2023-02-15T18:42:56, 1, null, null, null]
    -U[100001, user_001, 12345, 1, phone_aaa, 5000.05, 2023-02-15T16:40:56, 2023-02-15T18:42:56, 1, null, null, null]
    +U[100001, user_001, 12345, 1, phone_aaa, 5000.05, 2023-02-15T16:40:56, 2023-02-15T18:42:56, 1, 2001, 1, 2023-02-15T17:40:56]
    +U[100004, user_001, 12347, 4, phone_ddd, 2000.02, 2023-02-15T13:40:56, 2023-02-15T18:42:56, 1, 2004, 0, 2023-02-15T17:40:56]
    +U[100006, user_001, 12348, 1, phone_aaa, 1000.01, 2023-02-15T11:40:56, 2023-02-15T18:42:56, 1, 2006, 0, 2023-02-15T18:40:56]

Batch-mode profiling

A análise exploratória em modo batch não grava dados em uma tabela de resultado. Em vez disso, recupera o estado atual dos dados e permite visualizar os resultados diretamente por meio de depuração.

Na página Data Development>ETL, crie um SQL Stream Job, copie o código a seguir no editor SQL e clique em Debug. Para mais informações, consulte Depuração de Jobs.

O resultado da depuração na interface de desenvolvimento de jobs do Flink é mostrado abaixo.

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 the execution efficiency of batch jobs.

Após a depuração, a interface de desenvolvimento de jobs do Flink retorna dois registros de pedidos que correspondem aos critérios de filtro: os valores de order_id são 100004 e 100001, ambos para order_user_id user_001. Seus valores de order_fee são 2000.02 e 5000.05, e seus valores de order_create_time são posteriores a 2023-02-15 12:00:00.

Etapa 3: Usar o data warehouse em tempo real

Após construir o Streaming Warehouse em camadas com o catálogo do Flink, utilize o data warehouse para consultas pontuais, análises OLAP e relatórios em tempo real.

Consultas pontuais

Consulte as tabelas de métricas agregadas na camada DWS com base na chave primária, suportando milhões de RPS.

Veja a seguir um exemplo de código para consultar o valor de consumo de um usuário específico em uma data específica na página de desenvolvimento do HoloWeb.

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

No resultado da consulta, o valor do campo de montante de consumo (paied_buy_fee_sum) é 8000.08.

Análise OLAP

Realize análises OLAP na tabela ampla da camada DWD.

Veja a seguir um exemplo de código para consultar os detalhes de pedidos de um cliente específico em uma plataforma de pagamento específica em fevereiro de 2023 na página de desenvolvimento do HoloWeb.

-- 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;

Após a execução da consulta, a tabela de resultados exibe os registros de detalhes de pedidos que correspondem aos critérios de filtro, mostrando 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.

Relatórios em tempo real

Exiba relatórios em tempo real a partir da tabela ampla da camada DWD. O armazenamento híbrido linha-coluna e as tabelas orientadas a colunas do Hologres oferecem forte desempenho de análise OLAP com tempos de resposta no nível de segundos.

Veja a seguir um exemplo de código para consultar o número total e o valor total de pedidos para cada categoria em fevereiro de 2023 na página de desenvolvimento do HoloWeb.

-- 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 o SQL, a aba Results exibe quatro colunas de dados em uma tabela: order_create_date, order_product_catalog_name, count e sum. O resultado de exemplo para a data 20230215 mostra as contagens de pedidos (2, 1, 1, 2 e 2) e os valores totais (6000.06, 4000.04, 3000.03, 4000.04 e 7000.03) para cinco categorias de product de phone_aaa a phone_eee.

Referências