Tous les produits
Search
Centre de documentation

Hologres:Construire un entrepôt de données en temps réel avec Flink et Hologres

Dernière mise à jour :Aug 11, 2026

Combinez le traitement en temps réel de Realtime Compute for Apache Flink avec les fonctionnalités de Hologres, telles que Binlog, le stockage hybride ligne-colonne et l'isolation stricte des ressources, pour construire un entrepôt de données en temps réel évolutif capable de gérer des volumes de données croissants et les exigences métier en temps réel.

Contexte

La demande croissante de fraîcheur des données pousse les entreprises à dépasser le traitement par lots hors ligne traditionnel au profit du traitement, du stockage et de l'analyse des données en temps réel. Alors que l'entreposage de données hors ligne suit une méthodologie bien définie avec un traitement en couches (ODS > DWD > DWS > ADS) via des tâches planifiées, un cadre comparable pour l'entreposage de données en temps réel est encore en cours d'émergence. Cette solution applique le concept de Streaming Warehouse pour créer un flux efficace de données en temps réel entre les couches et répondre aux défis de la stratification des données en temps réel.

Scénario

En prenant l'exemple d'une plateforme de commerce électronique, cette rubrique montre comment construire un entrepôt de données en temps réel en intégrant Flink à Hologres. La configuration résultante permet le traitement et le nettoyage des données en temps réel, prend en charge les requêtes des applications en amont et établit la stratification et la réutilisation des données pour des scénarios tels que les tableaux de bord transactionnels, l'analyse comportementale, la création de profils utilisateurs et les recommandations personnalisées.

Architecture

image
  1. Construisez la couche ODS : ingérez les données des bases de données métier en temps réel.

    MySQL contient trois tables métier : orders (table des commandes), orders_pay (table des paiements de commandes) et product_catalog (table de dictionnaire des catégories de produits). Flink synchronise ces trois tables vers Hologres en temps réel pour former la couche ODS.

  2. Construisez la couche DWD : créez une table large en temps réel.

    Flink joint les tables orders, product_catalog et orders_pay en temps réel pour créer une table large au niveau de la couche DWD.

  3. Construisez la couche DWS : calculez les métriques en temps réel.

    Flink consomme le Binlog de la table large via un processus piloté par les événements, en agrégeant les données pour créer des tables de métriques pour les utilisateurs et les boutiques au niveau de la couche DWS.

  4. Servez les requêtes des applications via Hologres.

    • Les applications peuvent interroger les tables de métriques agrégées au niveau de la couche DWS, prenant en charge des millions de requêtes par seconde (RPS).

    • Les applications peuvent effectuer une analyse OLAP sur la table large DWD ou afficher des rapports en temps réel basés sur ses données, avec des temps de réponse de quelques secondes.

Avantages et fonctionnalités principales

Principaux avantages :

  • Mises à jour efficaces et requêtes immédiates : Hologres prend en charge des mises à jour et des corrections efficaces, ainsi que la cohérence lecture-après-écriture pour les données de chaque couche. Cela remédie à une limitation majeure des entrepôts de données en temps réel traditionnels, où il est difficile d'interroger, de mettre à jour et de corriger les données des couches intermédiaires.

  • Stratification et réutilisation des données : Toutes les couches de données dans Hologres peuvent être exposées indépendamment aux services externes, permettant une stratification et une réutilisation efficaces des données.

  • Architecture simplifiée et efficacité améliorée : La construction d'un pipeline ETL en temps réel avec Flink SQL et le stockage des données des couches ODS, DWD et DWS dans Hologres simplifient l'architecture et améliorent l'efficacité du traitement des données.

Cette solution s'appuie sur trois fonctionnalités principales de Hologres.

Fonctionnalité principale

Description

Binlog

Le Binlog de Hologres pilote Flink pour effectuer des calculs en temps réel et sert de source en amont pour le traitement de flux.

Stockage hybride ligne-colonne

Hologres prend en charge un format de stockage hybride ligne-colonne. Une seule table stocke les données à la fois dans des formats orientés ligne et orientés colonne avec une forte cohérence. Les tables intermédiaires peuvent servir de tables source pour Flink, de tables de dimension pour les requêtes ponctuelles et les jointures, et peuvent également être interrogées par d'autres applications telles que les services OLAP et en ligne.

Isolation stricte des ressources

Une charge élevée sur une instance Hologres peut affecter les performances des requêtes ponctuelles sur les couches intermédiaires. Hologres prend en charge une isolation stricte des ressources grâce au déploiement avec séparation lecture/écriture pour les instances principales et secondaires (stockage partagé) ou à l'architecture d'instance de virtual warehouse. Cela garantit que les jobs Flink qui extraient les données Binlog de Hologres n'affectent pas les services en ligne.

Prérequis

  • Seules les instances exclusives de Hologres prennent en charge cette solution d'entrepôt de données en temps réel.

  • Les instances Realtime Compute for Apache Flink, RDS MySQL et Hologres doivent se trouver dans le même VPC. Si ce n'est pas le cas, vous devez d'abord connecter les VPC ou utiliser des endpoints publics pour l'accès. Pour plus d'informations, consultez les rubriques Comment accéder à d'autres services à travers les VPC ? et Comment accéder à Internet ?.

  • Assurez-vous que tout utilisateur RAM ou rôle RAM utilisé pour l'accès dispose des autorisations requises pour les ressources Realtime Compute for Apache Flink, Hologres et RDS MySQL.

Étape 1 : Préparer les ressources

Créer une instance RDS MySQL et préparer une source de données

  1. Créez une instance RDS MySQL. Pour plus d'informations, consultez la rubrique Créer une instance ApsaraDB RDS for MySQL.

    L'instance RDS MySQL doit se trouver dans le même VPC que l'espace de travail Flink et l'instance Hologres.

  2. Créez une base de données et un compte.

    Sur l'instance cible, créez une base de données nommée order_dw et un compte standard disposant des autorisations de lecture et d'écriture pour cette base de données. Pour plus de détails, consultez les rubriques Créer une base de données et Créer un compte.

  3. Préparez la source de données MySQL CDC.

    1. Sur la page des détails de l'instance, cliquez sur Log On to Database.

    2. Dans la boîte de dialogue Connect to Instance, saisissez le nom d'utilisateur et le mot de passe du compte que vous avez créé, puis cliquez sur Sign in.

    3. Une fois connecté, double-cliquez sur la base de données order_dw sur la page de l'instance de base de données pour basculer vers celle-ci.

    4. Dans la zone SQL Console, rédigez les instructions DDL pour créer les trois tables métier ainsi que les instructions d'insertion des données.

      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. Cliquez sur Upload, puis cliquez sur Execute.

Créer une instance Hologres et un virtual warehouse

  1. Créez une instance Hologres exclusive. Pour plus d'informations, consultez la rubrique Acheter une instance Hologres.

    L'instance Hologres doit se trouver dans le même VPC que l'instance RDS MySQL. Pour bénéficier de la capacité d'isolation stricte des ressources de Hologres via la séparation lecture/écriture, sélectionnez Virtual Warehouse comme type d'instance et définissez Reserved Computing Resources of Virtual Warehouse sur 64. Cela vous permet de créer un nouveau virtual warehouse.

  2. Après vous être connecté à l'instance, créez une base de données et accordez des autorisations.

    Créez une base de données nommée order_dw (le modèle d'autorisation simple doit être activé) et accordez des autorisations d'administrateur à l'utilisateur. Pour plus d'informations sur la création d'une base de données et l'octroi d'autorisations, consultez la rubrique DB Management.

    Remarque
    • Si un compte n'apparaît pas dans la liste déroulante User Account, il n'a pas été ajouté à l'instance. Accédez à la page Users pour ajouter l'utilisateur en tant que SuperUser.

    • Dans Hologres V2.0 et versions ultérieures, l'extension Binlog est activée par défaut et ne nécessite pas d'exécution manuelle.

  3. Créez un nouveau virtual warehouse.

    Vous pouvez utiliser différents virtual warehouses pour obtenir une isolation des ressources. Utilisez le virtual warehouse initial init_warehouse pour l'écriture des données, et le virtual warehouse read_warehouse_1 pour servir les requêtes.

    Les ressources de calcul réservées sont entièrement allouées au virtual warehouse initial init_warehouse. Vous devez d'abord réduire ses ressources avant d'en créer un nouveau. Pour plus d'informations, consultez la rubrique Créer une nouvelle instance de virtual warehouse.

    1. Accédez à Security Center > Virtual Warehouse Management, et confirmez que le nom de l'instance est correct.

    2. Dans la ligne du virtual warehouse existant init_warehouse, cliquez sur Modify Configuration dans la colonne Actions. Réduisez les ressources et cliquez sur OK.

    3. Cliquez sur Create Virtual Warehouse, créez un nouveau virtual warehouse nommé read_warehouse_1, puis cliquez sur OK.

Créer un espace de travail Flink et des catalogues

  1. Créez un espace de travail Flink. Pour plus d'informations, consultez la rubrique Activer Realtime Compute for Apache Flink.

    L'espace de travail Flink doit se trouver dans le même VPC que les instances RDS MySQL et Hologres.

  2. Connectez-vous à la console Realtime Compute for Apache Flink et, sur la ligne de l'espace de travail cible, cliquez sur Console dans la colonne Actions.

  3. Créez un cluster de session pour fournir un environnement d'exécution permettant de créer des catalogues et des scripts de requête. Pour plus d'informations, reportez-vous à l'étape Étape 1 : Créer un cluster de session.

  4. Créez un catalogue Hologres.

    Dans l'onglet Script de la page Development > Scripts, copiez le code suivant dans l'éditeur de script. Modifiez les valeurs des paramètres cibles, sélectionnez l'extrait de code souhaité, puis cliquez sur Run. Dans le coin inférieur droit, assurez-vous que le cluster de session que vous avez créé est bien sélectionné comme environnement d'exécution.

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

    Vous devez modifier les valeurs des paramètres suivants en fonction des informations réelles de votre service Hologres.

    Parameter

    Description

    Notes

    endpoint

    Le endpoint de l'instance Hologres.

    Obtenez le nom de domaine pour le type de réseau VPC spécifié sur la page des détails de l'instance Hologres. Pour plus d'informations sur les noms de domaine, consultez la section Endpoints.

    username

    Sélectionnez l'une des options suivantes :

    • Le nom d'utilisateur du compte personnalisé est au format BASIC$<user_name>.

    • L'AccessKey ID de votre compte Alibaba Cloud ou d'un utilisateur RAM.

    • L'utilisateur configuré ici doit disposer d'un accès à la base de données Hologres correspondante. Pour plus d'informations sur les autorisations de base de données Hologres et la gestion des utilisateurs, consultez les rubriques Hologres permission model et User management.

    • Cet exemple utilise un compte personnalisé nommé BASIC$flinktest et définit son mot de passe à l'aide d'une variable de projet nommée holosecrect afin d'éviter les risques de sécurité liés au stockage en clair. Pour plus d'informations, consultez la rubrique Project variables.

    password

    • Le mot de passe du compte personnalisé.

    • L'AccessKey Secret de votre compte Alibaba Cloud ou d'un utilisateur RAM.

    Remarque

    Lors de la création d'un catalogue, vous pouvez définir des paramètres WITH par défaut pour les tables source, de dimension et de résultat. Vous pouvez également configurer des propriétés par défaut pour la création de tables physiques Hologres, telles que les paramètres commençant par table_property. Pour plus d'informations, consultez les rubriques Gérer les catalogues Hologres et Paramètres WITH du connecteur Hologres.

  5. Créez un catalogue MySQL.

    Copiez le code suivant dans l'éditeur Script. Modifiez les valeurs des paramètres cibles, sélectionnez l'extrait de code souhaité, puis cliquez sur Run à gauche de la ligne de code. Dans le coin inférieur droit, vérifiez que le cluster de session que vous avez créé est bien sélectionné comme environnement d'exécution.

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

    Vous devez adapter les valeurs des paramètres suivants aux informations réelles de votre service MySQL.

    Parameter

    Description

    hostname

    L'adresse IP ou le nom d'hôte de la base de données MySQL. Vous pouvez obtenir l'adresse interne en cliquant sur View Details dans la zone Network Type de la page des informations de base de la base de données.

    port

    Le numéro de port du service de base de données MySQL. La valeur par défaut est 3306.

    username

    Le nom d'utilisateur du service de base de données MySQL.

    password

    Le mot de passe du service de base de données MySQL.

    Cet exemple utilise une variable nommée mysql_pw pour la valeur du mot de passe afin d'éviter les risques tels que le stockage en clair. Pour plus d'informations, consultez la rubrique Gérer les variables.

Étape 2 : Construire l'entrepôt de données en temps réel

Construire la couche ODS : Ingestion de données en temps réel

Construisez la couche ODS en une seule étape à l'aide de l'instruction CREATE DATABASE AS (CDAS) du catalogue. La couche ODS sert généralement de pilote d'événements pour les jobs de streaming plutôt que pour les requêtes OLAP ou les requêtes ponctuelles clé-valeur ; l'activation du Binlog suffit donc. Le connecteur Hologres prend en charge un mode complet et incrémentiel qui lit d'abord toutes les données existantes, puis consomme le Binlog de manière incrémentielle.

  1. Créez un job de synchronisation CDAS nommé ODS.

    1. Sur la page Data Development>ETL, créez un nouveau job de streaming SQL nommé ODS et copiez le code suivant dans l'éditeur 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.
      Remarque
      • Par défaut, cet exemple synchronise les données vers le schéma Public de la base de données order_dw. Vous pouvez également synchroniser les données vers un schéma spécifique de la base de données Hologres cible. Pour plus d'informations, consultez la section Utiliser comme catalogue côté cible pour CDAS. Après avoir spécifié un schéma, le format du nom de table lors de l'utilisation du catalogue change également. Pour plus d'informations, consultez la section Utiliser un catalogue Hologres.

      • Les modifications de schéma apportées à une table source sont propagées à la table de résultat uniquement après qu'une opération DML ultérieure (INSERT, UPDATE ou DELETE) a été effectuée sur la table source.

    2. Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.

    3. Dans le volet de navigation de gauche, cliquez sur O&M > Deployments. Sur la ligne du job ODS que vous venez de déployer, cliquez sur Start dans la colonne Actions. Sélectionnez Start Without State puis cliquez sur Start.

  2. Chargez les données dans l'entrepôt virtuel.

    Un groupe de tables constitue le support de données dans Hologres. Lorsque vous utilisez l'entrepôt virtuel read_warehouse_1 pour interroger les données d'un groupe de tables de la base de données order_dw, tel que order_dw_tg_default (pour les étapes de création, consultez la rubrique Gestion des groupes de tables), vous chargez order_dw_tg_default pour read_warehouse_1. Cela vous permet d'utiliser l'entrepôt virtuel init_warehouse pour écrire des données et l'entrepôt virtuel read_warehouse_1 pour effectuer des requêtes de service.

    Sur la page de développement HoloWeb, cliquez sur SQL Editor, confirmez le nom de l'instance et le nom de la base de données, puis exécutez les commandes suivantes. Pour plus d'informations, consultez l'étape Créer une nouvelle instance d'entrepôt virtuel. Après le chargement, vous pouvez constater que read_warehouse_1 a chargé les données du groupe de tables 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. Dans le coin supérieur droit, basculez l'entrepôt virtuel vers read_warehouse_1. Les requêtes et analyses suivantes utiliseront l'entrepôt virtuel read_warehouse_1.

    La liste déroulante des entrepôts virtuels dans le coin supérieur droit affiche read_warehouse_1. L'éditeur montre l'instruction de chargement du groupe de tables exécutée CALL hg_table_group_load_to_warehouse ('order_dw.order_dw_tg_default', 'read_warehouse_1', 1); ainsi que l'instruction de requête select * from hologres.hg_warehouse_table_groups;.

  4. Sur la page SQL Editor, exécutez les commandes suivantes pour afficher les données synchronisées depuis MySQL vers les trois tables 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;

    Après avoir exécuté la troisième requête, l'onglet Result[3] affiche cinq lignes dans la table product_catalog, comprenant deux colonnes : product_id (avec des valeurs allant de 1 à 5) et catalog_name (avec les valeurs phone_aaa, phone_bbb, phone_ccc, phone_ddd et phone_eee).

Construire la couche DWD : Créer une table large en temps réel

La couche DWD tire parti de la capacité de mise à jour partielle des colonnes du connecteur Hologres, permettant aux instructions DML INSERT d'exprimer une sémantique de mise à jour partielle des colonnes. Ce processus s'appuie sur des requêtes ponctuelles haute performance contre les tables de dimension grâce au stockage en ligne et au stockage hybride ligne-colonne de Hologres, tandis qu'une isolation forte des ressources garantit que les jobs d'écriture, de lecture et d'analyse ne s'interfèrent pas mutuellement.

  1. Utilisez la fonctionnalité de catalogue Flink pour créer la table large dwd_orders de la couche DWD dans Hologres.

    Dans l'onglet Script de la page Development > Scripts, copiez le code suivant dans l'éditeur de script, sélectionnez l'extrait, puis cliquez sur Run à gauche de la ligne de code.

    -- 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. Consommez le Binlog des tables de la couche ODS orders et orders_pay en temps réel.

    Sur la page Data Development>ETL, créez un job de streaming SQL nommé DWD, copiez le code suivant dans l'éditeur SQL, puis Deploy et Start le job. Ce job SQL joint la table orders à la table de dimension product_catalog, écrit le résultat final dans la table dwd_orders et effectue un enrichissement des données en temps réel.

    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. Consultez les données de la table large dwd_orders.

    Sur la page de développement HoloWeb, connectez-vous à l'instance Hologres et connectez-vous à la base de données cible. Ensuite, exécutez la commande suivante dans l'éditeur SQL.

    SELECT * FROM dwd_orders;

    La table large dwd_orders comprend les champs suivants : 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 et pay_create_time. La requête renvoie 7 enregistrements de commande.

Construction de la couche DWS : calcul des métriques en temps réel

  1. Utilisez la fonctionnalité de catalogue Flink pour créer les tables d'agrégation de la couche DWS dws_users et dws_shops dans Hologres.

    Dans l'onglet Script de la page Development > Scripts, copiez le code suivant dans l'éditeur de script, sélectionnez l'extrait, puis cliquez sur Run à gauche de la ligne de code.

    -- 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. Consommez les données de la table large DWD dw.order_dw.dwd_orders en temps réel, effectuez les agrégations dans Flink et écrivez les résultats finaux dans les tables DWS de Hologres.

    Sur la page Data Development>ETL, créez une nouvelle tâche de streaming SQL nommée DWS, copiez le code suivant dans l'éditeur SQL, puis Deploy et Start la tâche.

    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. Consultez les résultats d'agrégation au niveau de la couche DWS. Les résultats sont mis à jour en temps réel à mesure que les données amont évoluent.

    1. Dans la console Hologres, consultez les données avant modification.

      dws_users

      SELECT * FROM dws_users;

      Le résultat de la requête contient trois colonnes : user_id (par exemple, user_001, user_002, user_003), ds (par exemple, 20230215) et paied_buy_fee_sum (par exemple, 8000.08, 5000.05). La colonne user_id sert de clé d'association.

      dws_shops

      SELECT * FROM dws_shops;

      La requête renvoie 4 enregistrements avec trois colonnes : shop_id, ds et paied_buy_fee_sum. Dans les données d'exemple, shop_id varie de 12345 à 12348, toutes les valeurs ds sont 20230215 et les valeurs paied_buy_fee_sum sont respectivement 5000,05, 4000,04, 7000,07 et 2000,02.

    2. Dans la console RDS, insérez un nouvel enregistrement de données dans chacune des tables orders et orders_pay de la base de données 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. Dans la console Hologres, consultez les données après modification.

      dwd_orders

      SELECT * FROM dwd_orders;

      Le résultat d'exécution affiche un total de 8 enregistrements de commande dans la table dwd_orders (order_id 100001–100008). Les champs incluent order_user_id, order_shop_id, order_product_id, order_product_catalog_name, order_fee, order_create_time et order_update_time. L'enregistrement nouvellement inséré pour order_id=100008 présente un order_fee de 6000,02.

      dws_users

      SELECT * FROM dws_users;

      La requête renvoie 3 enregistrements de la table dws_users avec les colonnes user_id, ds et paied_buy_fee_sum. Les données sont : user_001 / 20230215 / 8000,08, user_002 / 20230215 / 5000,05 et user_003 / 20230215 / 11000,07. Le montant agrégé pour user_003 est le plus élevé, soit 11000,07.

      dws_shops

      SELECT * FROM dws_shops;

      La requête renvoie quatre lignes avec trois colonnes : shop_id, ds et paied_buy_fee_sum. Les valeurs shop_id sont 12345, 12346, 12347 et 12348 ; toutes les valeurs ds sont 20230215 ; et les valeurs paied_buy_fee_sum sont respectivement 11000,07, 4000,04, 7000,07 et 2000,02.

Profilage des données

Avec Binlog activé, vous pouvez inspecter directement les modifications de données. La persistance des données à chaque couche simplifie le profilage des données ad hoc et la vérification des résultats.

Stream-mode profiling

Vous pouvez utiliser le connecteur Print pour vérifier si les messages envoyés vers d'autres tables de résultats correspondent aux attentes.

  1. Créez et démarrez une tâche de profilage en flux.

    Sur la page Data Development>ETL, créez une tâche de streaming SQL nommée Data-exploration, copiez le code suivant dans l'éditeur SQL, puis Deploy et Start la tâche.

    -- 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. Consultez les résultats du profilage des données.

    Sur la page de détails O&M > Deployments, cliquez sur le nom de la tâche cible. Dans l'onglet Logs, cliquez sur l'onglet Operational Logs à gauche. Ensuite, cliquez sur l'onglet Running Task Managers et cliquez sur un Path, ID. Sur la page Stdout, recherchez les informations de journal liées à 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

Le profilage en mode batch n'écrit pas les données dans une table de résultats. Il récupère l'état actuel des données et vous permet de visualiser les résultats directement via le débogage.

Sur la page Data Development>ETL, créez une tâche SQL Stream, copiez le code suivant dans l'éditeur SQL et cliquez sur Debug. Pour plus d'informations, consultez Job Debugging.

Le résultat du débogage dans l'interface de développement de tâches Flink est illustré ci-dessous.

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.

Après le débogage, l'interface de développement de tâches Flink renvoie deux enregistrements de commande correspondant aux critères de filtrage : les valeurs order_id sont 100004 et 100001, toutes deux pour order_user_id user_001. Leurs valeurs order_fee sont 2000,02 et 5000,05, et leurs valeurs order_create_time sont postérieures au 15 février 2023 à 12:00:00.

Étape 3 : Utilisation de l'entrepôt de données en temps réel

Une fois l'entrepôt de données en streaming stratifié construit avec le catalogue Flink, vous pouvez exploiter l'entrepôt de données pour des requêtes ponctuelles, des analyses OLAP et des rapports en temps réel.

Requêtes ponctuelles

Interrogez les tables de métriques agrégées au niveau de la couche DWS en vous basant sur la clé primaire, avec prise en charge de millions de RPS.

Voici un exemple de code pour interroger le montant des dépenses d'un utilisateur spécifique à une date donnée sur la page de développement HoloWeb.

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

Dans le résultat de la requête, la valeur du champ de montant des dépenses (paied_buy_fee_sum) est 8000,08.

Analyse OLAP

Effectuez une analyse OLAP sur la table large de la couche DWD.

Voici un exemple de code pour interroger les détails des commandes d'un client spécifique sur une plateforme de paiement donnée en février 2023 sur la page de développement 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;

Une fois la requête exécutée, la table de résultats affiche les enregistrements de détails de commande correspondant aux critères de filtrage, présentant des champs tels que order_id, order_user_id, order_shop_id, order_product_id, order_product_catalog_name, order_fee, order_create_time et order_update_time.

Rapports en temps réel

Affichez des rapports en temps réel à partir de la table large de la couche DWD. Le stockage hybride ligne-colonne de Hologres et les tables orientées colonnes offrent de solides performances d'analyse OLAP avec des temps de réponse de l'ordre de la seconde.

Voici un exemple de code pour interroger le nombre total et le montant total des commandes pour chaque catégorie en février 2023 sur la page de développement 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;

Après avoir exécuté le code SQL, l'onglet Results affiche quatre colonnes de données dans un tableau : order_create_date, order_product_catalog_name, count et sum. L'exemple de résultat pour la date 20230215 montre les nombres de commandes (2, 1, 1, 2 et 2) et les montants totaux (6000,06, 4000,04, 3000,03, 4000,04 et 7000,03) pour cinq catégories de produits, de phone_aaa à phone_eee.

Références