Ce guide explique comment construire un entrepôt de données temps réel à l'aide de Realtime Compute for Apache Flink et Hologres. Cette solution associe la puissance du traitement de flux de Flink aux fonctionnalités uniques de Hologres — telles que le binary logging, le stockage hybride ligne-colonne et une forte isolation des ressources — afin de gérer des volumes de données croissants et de répondre aux exigences métier en temps réel.
Contexte
Avec la digitalisation croissante des entreprises, le besoin de données fraîches s'intensifie rapidement. Les entrepôts de données hors ligne traditionnels, conçus pour le traitement par lots de gros volumes, ne suffisent plus. De nombreux scénarios métier modernes exigent désormais un traitement, un stockage et une analyse des données en temps réel. Bien que la méthodologie de construction d'entrepôts de données hors ligne reposant sur des architectures en couches (ODS, DWD, DWS) soit bien établie, il manquait jusqu'à présent un cadre clair pour leurs équivalents temps réel. L'utilisation d'un entrepôt de données temps réel permet d'assurer un flux de données efficace et immédiat entre chaque couche de données.
Cas d'utilisation
Ce guide prend l'exemple d'une plateforme de e-commerce pour illustrer la construction d'un entrepôt de données temps réel intégrant Flink et Hologres. Cette approche permet de traiter et de nettoyer les données en temps réel, de fournir aux applications en aval des données structurées et réutilisables, et de prendre en charge divers scénarios métier tels que les tableaux de bord temps réel (suivi des transactions, analyse comportementale, profilage utilisateur) ou les recommandations personnalisées.
Architecture de la solution
-
Construire la couche ODS (Operational Data Store) : ingérer les données de la base de données métier en temps réel.
Flink synchronise trois tables métier MySQL —
orders(table des commandes),orders_pay(table des paiements) etproduct_catalog(dictionnaire des catégories de produits) — vers Hologres en temps réel. Ces tables constituent la couche ODS. -
Construire la couche DWD (Data Warehouse Detail) : créer une table large temps réel.
Flink joint les tables ODS en temps réel pour produire une table large destinée à la couche DWD.
-
Construire la couche DWS (Data Warehouse Service) : calculer les métriques temps réel.
Flink consomme les modifications du binary logging de la table large selon un modèle événementiel, puis agrège les métriques dans des tables spécifiques aux utilisateurs et aux boutiques pour la couche DWS.
-
Répondre aux requêtes applicatives via Hologres.
Interroger les tables de métriques agrégées de la couche DWS, en supportant des millions de requêtes par seconde (RPS).
Exécuter des requêtes OLAP sur la table large DWD ou afficher des rapports temps réel basés sur ses données, avec des temps de réponse de l'ordre de la seconde.
Avantages et capacités principales
Cette solution offre les avantages suivants :
Mises à jour efficaces et requêtes immédiates : Hologres permet des mises à jour et corrections performantes, ainsi qu'un accès immédiat par requête pour chaque couche de données. Cela résout les difficultés courantes des entrepôts de données temps réel traditionnels, où les données intermédiaires sont souvent complexes à interroger, mettre à jour ou corriger.
Structuration et réutilisation des données : Chaque couche de données dans Hologres peut servir indépendamment les applications externes. Cela favorise une réutilisation efficace des données et permet de constituer un entrepôt structuré et modulaire.
Architecture simplifiée et efficacité accrue : L'utilisation de Flink SQL pour construire le pipeline ETL temps réel, combinée au stockage de toutes les couches de données (ODS, DWD et DWS) dans Hologres, simplifie l'architecture globale et améliore l'efficacité du traitement des données.
Cette solution repose sur trois capacités fondamentales de Hologres, détaillées dans le tableau suivant.
|
Capacité principale |
Description |
|
Hologres fournit le binary logging, qui permet à Flink de lire les modifications de données en temps réel. Ainsi, Hologres agit comme source de streaming pour les jobs Flink. |
|
|
Hologres prend en charge un format de stockage hybride où une même table conserve les données à la fois en mode orienté ligne et en mode orienté colonne, avec une forte cohérence. Une table intermédiaire peut ainsi servir de source Flink, de table de dimension pour les requêtes ponctuelles et les jointures temporelles, ou encore de source de données pour d'autres applications comme les requêtes OLAP ou les services en ligne. |
|
|
Forte isolation des ressources |
Une charge élevée sur une instance Hologres peut dégrader les performances des requêtes ponctuelles sur les couches de données intermédiaires. Hologres assure une forte isolation des ressources grâce à la séparation lecture/écriture entre instances primaire et secondaire (stockage partagé) ou à l'architecture virtual warehouse. Cela garantit que l'ingestion de données par Flink depuis les binary logs n'impacte pas les services en ligne. |
Remarques d'utilisation
Cette solution d'entrepôt de données temps réel est prise en charge uniquement sur les instances Hologres dédiées.
Votre espace de travail Realtime Compute for Apache Flink, votre instance ApsaraDB RDS for MySQL et votre instance Hologres doivent se trouver dans le même VPC. S'ils sont situés dans des VPC différents, vous devez d'abord les connecter ou utiliser des endpoints publics. Pour plus d'informations, consultez Comment accéder à d'autres services entre VPC ? et Comment accéder à Internet ?.
Si vous utilisez un utilisateur RAM ou un rôle RAM pour accéder aux ressources Realtime Compute for Apache Flink, Hologres et ApsaraDB RDS for MySQL, assurez-vous qu'il dispose des permissions nécessaires.
Étape 1 : Préparer l'environnement
Créer une instance RDS for MySQL et préparer les données
-
Créez une instance ApsaraDB RDS for MySQL. Pour plus d'informations, consultez Créer une instance ApsaraDB RDS for MySQL.
L'instance ApsaraDB RDS for MySQL doit se trouver dans le même VPC que votre espace de travail Flink et votre instance Hologres.
-
Créez une base de données et un compte.
Pour l'instance cible, créez une base de données nommée
order_dwainsi qu'un compte standard disposant des permissions de lecture et d'écriture sur cette base. Pour plus d'informations, consultez Créer une base de données et Créer un compte. -
Préparez la source de données MySQL CDC.
Sur la page de détails de l'instance, cliquez sur Log On to Database.
Sur la page de connexion, saisissez le nom d'utilisateur et le mot de passe du compte de base de données créé, puis cliquez sur Log On.
Une fois connecté, double-cliquez sur la base de données
order_dwpour basculer dessus.-
Dans la console SQL, saisissez les instructions DDL suivantes pour créer les tables métier, ainsi que les instructions INSERT pour les alimenter.
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');
Cliquez sur Execute, puis sur Direct Execution.
Créer une instance Hologres et des compute groups
-
Achetez une instance Hologres dédiée. Pour plus d'informations, consultez Acheter une instance Hologres.
L'instance Hologres doit se trouver dans le même VPC que l'instance ApsaraDB RDS for MySQL. Pour bénéficier d'une forte isolation des ressources via la séparation lecture/écriture, cet exemple utilise le type d'instance Virtual Warehouse et définit la Reserved Compute Resource à 64, ce qui permet de créer des compute groups supplémentaires.
-
Après vous être connecté à l'instance, créez une base de données et attribuez les permissions.
Créez une base de données nommée order_dw (avec le modèle de permissions simple activé) et accordez les privilèges administrateur à l'utilisateur. Pour plus de détails sur la gestion des bases de données et les autorisations, consultez Gérer les bases de données.
RemarqueSi le compte n'apparaît pas dans la liste déroulante User, cela signifie qu'il n'a pas été ajouté à l'instance. Accédez à la page User Management et ajoutez l'utilisateur en tant que SuperUser.
À partir de Hologres V2.0, l'extension binary logging est activée par défaut. Aucune activation manuelle n'est nécessaire.
-
Créez un nouveau compute group.
Différents compute groups permettent d'isoler les ressources. Utilisez le compute group initial
init_warehousepour les écritures de données et le compute groupread_warehouse_1pour répondre aux requêtes.Par défaut, toutes les ressources de calcul réservées sont allouées au compute group initial
init_warehouse. Vous devez d'abord réduire ses ressources avant de pouvoir créer un nouveau compute group. Pour plus d'informations, consultez Créer une instance de compute group.Accédez à et confirmez le nom de l'instance.
Dans la ligne correspondant au compute group
init_warehouse, cliquez sur Modify Configuration dans la colonne Actions. Réduisez les ressources allouées, puis cliquez sur OK.Cliquez sur Create Compute Group, créez un nouveau compute group nommé
read_warehouse_1, puis cliquez sur OK.
Créer un espace de travail Flink et des catalogs
-
Créez un espace de travail Flink. Pour plus d'informations, consultez Activer Realtime Compute for Apache Flink.
L'espace de travail Flink doit se trouver dans le même VPC que les instances ApsaraDB RDS for MySQL et Hologres.
Connectez-vous à la console Realtime Compute for Apache Flink et cliquez sur Console dans la colonne Actions de votre espace de travail.
Créez un cluster de session pour fournir un environnement d'exécution permettant de créer des catalogs et d'exécuter des scripts. Pour plus d'informations, consultez Étape 1 : Créer un cluster de session.
-
Créez un catalog Hologres.
Sur la page , dans l'onglet Scripts, copiez le code suivant, remplacez les valeurs d'espace réservé, sélectionnez le code, puis cliquez sur Run. Le cluster de session créé sert alors d'environnement d'exécution.
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' );Modifiez les paramètres suivants avec les informations réelles de votre service Hologres.
Paramètre
Description
Notes
endpoint
Endpoint de votre instance Hologres.
Sur la page de détails de l'instance Hologres, récupérez le nom de domaine correspondant au VPC spécifié. Pour plus d'informations sur les noms de domaine, consultez Endpoints.
username
Choisissez l'une des options suivantes :
-
Le nom d'utilisateur d'un compte personnalisé doit respecter le format
BASIC$< user_name>. -
L'AccessKey ID de votre compte Alibaba Cloud ou de votre utilisateur RAM.
-
L'utilisateur configuré doit avoir accès à la base de données Hologres correspondante. Pour plus de détails, consultez Modèle de permissions Hologres et Gérer les utilisateurs.
-
Cet exemple utilise un compte personnalisé nommé
BASIC$flinktestet définit son mot de passe via une variable de projet nommée holosecrect, afin d'éviter les risques de sécurité liés au stockage des mots de passe en clair. Pour plus d'informations, consultez Variables de projet.
password
-
Mot de passe du compte personnalisé.
-
AccessKey secret de votre compte Alibaba Cloud ou de votre utilisateur RAM.
RemarqueLors de la création d'un catalog, vous pouvez définir des options WITH par défaut pour les tables sources, de dimension et de résultat. Vous pouvez également configurer des propriétés par défaut pour les tables physiques Hologres, telles que les paramètres commençant par
table_property. Pour plus d'informations, consultez Gérer les catalogs Hologres et Connecteur Hologres pour les entrepôts de données temps réel (paramètres WITH). -
-
Créez un catalog MySQL.
Copiez le code suivant dans l'onglet Scripts, modifiez les valeurs des paramètres, sélectionnez le code, puis cliquez sur Run. Le cluster de session créé sert alors d'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' );Modifiez les paramètres suivants avec les informations réelles de votre service MySQL.
Paramètre
Description
hostname
Adresse IP ou nom d'hôte de votre base de données MySQL. Sur la page d'informations de base de la base de données, cliquez sur View Connection Details dans la zone Network Type pour obtenir l'endpoint interne.
port
Numéro de port du service de base de données MySQL. La valeur par défaut est 3306.
username
Nom d'utilisateur du service de base de données MySQL.
password
Mot de passe du service de base de données MySQL.
Cet exemple utilise une variable nommée
mysql_pwpour le mot de passe afin d'éviter toute exposition en clair. Pour plus d'informations, consultez Gérer les variables.
Étape 2 : Construire l'entrepôt de données temps réel
Construire la couche ODS : Ingérer les données métier
Grâce à l'instruction CREATE DATABASE AS (CDAS) basée sur les catalogs, vous pouvez créer la couche ODS en une seule étape. Cette couche sert généralement de source d'événements pour les jobs de streaming plutôt que de cible directe pour des requêtes OLAP ou ponctuelles. L'activation du binary logging suffit donc pour cet usage. Le binary logging est une capacité centrale de Hologres. Le connecteur Hologres prend également en charge un mode complet + incrémental : il lit d'abord un snapshot complet, puis consomme les binary logs de manière incrémentale.
-
Créez le job de synchronisation CDAS ODS.
-
Sur la page , créez un nouveau brouillon de flux SQL nommé ODS et copiez le code suivant dans l'éditeur 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.RemarquePar 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 Utiliser un catalog Hologres comme destination dans une instruction CREATE DATABASE AS.... Une fois un schéma spécifié, le format du nom de table change lors de l'utilisation du catalog. Pour plus de détails, consultez Utiliser un catalog Hologres.Si le schéma d'une table source change, le schéma de la table résultante ne sera mis à jour qu'après une modification de données (suppression, insertion ou mise à jour) dans la table source.
Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.
Dans le volet de navigation de gauche, choisissez . Dans la ligne du job ODS que vous venez de déployer, cliquez sur Start dans la colonne Actions. Sélectionnez Initial Mode, puis cliquez sur Start.
-
-
Chargez les données dans le compute group.
Un table group est un conteneur de données dans Hologres. Lorsque vous utilisez le compute group read_warehouse_1 pour interroger des données d'un table group dans la base de données order_dw, tel que order_dw_tg_default (pour créer un table group, consultez Table Group Management), le table group order_dw_tg_default est chargé pour le compute group read_warehouse_1. Cela vous permet d'utiliser le compute group
init_warehousepour écrire des données et le compute groupread_warehouse_1pour les requêtes de service.Sur la page de développement HoloWeb, cliquez sur SQL Editor. Confirmez le nom de l'instance et celui de la base de données, puis exécutez les commandes suivantes. Pour plus d'informations, consultez Créer une instance de compute group. Après le chargement, vous pouvez constater que
read_warehouse_1a bien chargé les données du table grouporder_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; -
Dans le coin supérieur droit, basculez le compute group vers
read_warehouse_1. Les requêtes et analyses ultérieures utiliseront ce compute group.Dans le coin supérieur droit de la page HoloWeb, sélectionnez
read_warehouse_1dans la liste déroulante des compute groups. -
Sur la page SQL Editor, exécutez les commandes suivantes pour visualiser les données synchronisées depuis MySQL vers les trois tables 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;Le résultat de la requête sur la table
product_catalogcontient deux colonnes, product_id (1 à 5) et catalog_name (phone_aaa, phone_bbb, phone_ccc, phone_ddd, phone_eee), pour un total de 5 enregistrements. Cela confirme que les données ont été correctement synchronisées vers Hologres.
Construire la couche DWD : Créer une table large temps réel
Cette étape exploite la capacité de mise à jour partielle des colonnes du connecteur Hologres. Vous pouvez exprimer ces mises à jour partielles via des instructions DML INSERT. Le job interroge plusieurs tables de dimension au moyen de requêtes ponctuelles haute performance, rendues possibles par le stockage ligne et le stockage hybride ligne-colonne de Hologres. Grâce à la forte isolation des ressources, les charges d'écriture, de lecture et d'analyse n'interfèrent pas entre elles.
-
Utilisez la fonctionnalité de catalog Flink pour créer la table large DWD
dwd_ordersdans Hologres.Sur la page , copiez le code suivant dans l'onglet Scripts, sélectionnez le code, puis cliquez sur 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. ); -
Consommez en temps réel les modifications du binary logging des tables ODS
ordersetorders_pay.Sur la page , créez un nouveau brouillon de flux SQL nommé DWD. Copiez le code suivant dans l'éditeur SQL, puis Deploy et Start le job. Ce job SQL joint les tables
ordersetproduct_catalogà l'aide d'une jointure temporelle et écrit le résultat dans la tabledwd_orders, enrichissant ainsi les 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; -
Visualisez les données de la table large
dwd_orders.Connectez-vous à l'instance Hologres sur la page de développement HoloWeb, connectez-vous à la base de données cible, puis exécutez la commande suivante dans l'éditeur SQL.
SELECT * FROM dwd_orders;Après une exécution réussie, le résultat de la requête renvoie les données de la table large dwd_orders, incluant 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.
Construire la couche DWS : Calculer les métriques temps réel
-
Utilisez la fonctionnalité de catalog Flink pour créer les tables d'agrégation DWS
dws_usersetdws_shopsdans Hologres.Sur la page , copiez le code suivant dans l'onglet Scripts, sélectionnez le code, puis cliquez sur 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 ); -
Consommez la table large DWD
dw. order_dw. dwd_ordersen temps réel, effectuez les agrégations dans Flink, puis écrivez les résultats finaux dans les tables DWS de Hologres.Sur la page , créez un nouveau brouillon de flux SQL nommé DWS. Copiez le code suivant dans l'éditeur SQL, puis Deploy et Start le 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; -
Consultez les résultats agrégés dans la couche DWS. Ces résultats sont mis à jour en temps réel à mesure que les données en amont changent.
-
Visualisez les données dans la console Hologres avant la modification.
Table dws_users
SELECT * FROM dws_users;Après l'exécution de la requête, le résultat renvoie les données de la table dws_users, qui comprend trois colonnes : user_id, ds et paied_buy_fee_sum. Dans l'exemple de résultat, la colonne user_id contient les valeurs
user_001,user_002etuser_003; la colonne ds contient la valeur20230215; et la colonne paied_buy_fee_sum contient respectivement les valeurs8000.08,5000.05et5000.05. La colonne user_id identifie de manière unique chaque utilisateur.Table dws_shops
SELECT * FROM dws_shops;Le résultat de la requête montre que la table
dws_shopscontient trois colonnes : shop_id (ID de la boutique), ds (partition de date) et paied_buy_fee_sum (montant du paiement). Quatre lignes de données d'exemple sont renvoyées, confirmant que la table de la couche DWS a été construite avec succès. -
Dans la console RDS, insérez un nouvel enregistrement dans chacune des tables
ordersetorders_payde la base de donnéesorder_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'); -
Visualisez les données mises à jour dans la console Hologres.
Table dwd_orders
SELECT * FROM dwd_orders;Après l'exécution de la requête, huit enregistrements de commande sont renvoyés depuis la table
dwd_orders, incluant 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. Le huitième enregistrement (order_id=100008,user_003,phone_eee,6000.02) correspond aux données nouvellement écrites.Table dws_users
SELECT * FROM dws_users;La requête renvoie trois lignes composées de trois colonnes : user_id, ds et paied_buy_fee_sum :
user_001 / 20230215 / 8000.08,user_002 / 20230215 / 5000.05etuser_003 / 20230215 / 11000.07. Le montant total de paiement le plus élevé concerne l'utilisateur user_003 (11000.07).Table dws_shops
SELECT * FROM dws_shops;Une fois la requête exécutée, le résultat contient trois colonnes, shop_id, ds et paied_buy_fee_sum, avec quatre lignes de données affichant les montants (11000.07, 4000.04, 7000.07 et 2000.02) pour les boutiques 12345, 12346, 12347 et 12348 à la date du 20230215. Les colonnes shop_id et paied_buy_fee_sum constituent les métriques clés.
-
Profiler les données
Puisque le binary logging est activé, vous pouvez inspecter directement les modifications de données. Si vous devez effectuer une exploration ad hoc des données métier sur des résultats intermédiaires ou vérifier l'exactitude du calcul final, chaque couche de cette solution est persistée, ce qui facilite l'examen du processus intermédiaire.
Profilage en mode streaming
Vous pouvez utiliser le connecteur Print pour vérifier si les messages envoyés aux autres tables de résultats correspondent aux attentes.
-
Créez et démarrez un job de profilage de données en streaming.
Sur la page , créez un nouveau brouillon de flux SQL nommé Data-exploration. Copiez le code suivant dans l'éditeur SQL, puis Deploy et Start le 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'; -
Consultez les résultats du profilage des données.
Sur la page de détails , cliquez sur le nom du job cible. Dans l'onglet Logs, cliquez sur l'onglet Logs, puis cliquez sur un lien Path, ID sous Running Task Managers. Sur la page Stdout, recherchez les informations de journal liées à
user_001.La sortie du journal affiche des enregistrements de modification de données CDC préfixés par
+I(insertion),-U(avant mise à jour) et+U(après mise à jour), incluant des champs tels queorder_id,order_user_id,order_shop_id,order_feeetorder_create_time.
Profilage en mode batch
Le profilage en mode batch n'écrit pas de données dans une table de résultats. Il récupère plutôt l'état final des données à l'instant présent, ce qui permet de visualiser les résultats directement dans la sortie de débogage.
Sur la page , créez un brouillon de flux SQL, copiez le code suivant dans l'éditeur SQL, puis cliquez sur Debug. Pour plus d'informations, consultez Déboguer un job.
Le résultat du débogage sur la page de développement de jobs Flink est le suivant.
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.
Une fois le débogage terminé, le résultat de la requête renvoie deux enregistrements de commande répondant aux conditions de filtrage, incluant 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. Cela confirme que le résultat du profilage en mode batch est conforme aux attentes.
Étape 3 : Utiliser l'entrepôt de données temps réel
L'étape 2 a montré comment utiliser les catalogs Flink pour construire un entrepôt de données temps réel structuré en couches, basé sur Flink et Hologres. Les sections suivantes décrivent plusieurs scénarios d'application simples.
Requête ponctuelle
Interrogez les tables de métriques agrégées de la couche DWS par clé primaire, avec une capacité de plusieurs millions de RPS.
Sur la page de développement HoloWeb, exécutez le code SQL suivant pour interroger le montant de consommation d'un utilisateur spécifique à une date donnée.
-- holo sql
SELECT * FROM dws_users WHERE user_id ='user_001' AND ds = '20230215';
Le résultat de la requête renvoie trois colonnes : user_id, ds et paied_buy_fee_sum (montant de consommation). Le montant de consommation pour user_001 à la date du 20230215 est de 8000.08.
Requête OLAP
Exécutez des requêtes OLAP sur la table large de la couche DWD.
Sur la page de développement HoloWeb, exécutez le code SQL suivant pour interroger les détails des commandes d'un client spécifique sur une plateforme de paiement donnée en février 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;
Le résultat de la requête renvoie deux enregistrements de commande avec des champs tels que order_id, order_user_id, order_shop_id, order_product_id, order_fee, order_create_time et order_update_time. Les ID de commande de l'exemple sont 100006 et 100004.
Rapports temps réel
Générez des rapports temps réel à partir des données de la table large de la couche DWD. Le stockage hybride ligne-colonne et les tables orientées colonne de Hologres offrent d'excellentes capacités de requête OLAP, garantissant des temps de réponse de l'ordre de la seconde.
Sur la page de développement HoloWeb, exécutez le code SQL suivant pour interroger le nombre total de commandes et le montant total des commandes pour chaque catégorie de produits en février 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;
Après l'exécution de l'instruction SQL, l'onglet Result affiche un tableau composé de quatre colonnes : order_create_date, order_product_catalog_name, count et sum. Par exemple, les résultats pour la date du 20230215 indiquent que la catégorie phone_aaa comptait 2 commandes pour un total de 6000.06, la catégorie phone_bbb comptait 1 commande pour 4000.04, et ainsi de suite pour les cinq catégories.
Références
-
Tutoriels pour des scénarios connexes :
Pour plus d'informations sur la capacité de binary logging de Hologres, consultez S'abonner au Binlog Hologres.
Flink prend en charge plusieurs instructions INSERT INTO dans un seul job. Pour plus de détails sur la syntaxe, consultez Instruction INSERT INTO.
Realtime Compute for Apache Flink prend en charge un large éventail de connecteurs. Pour plus d'informations, consultez Connecteurs pris en charge.