Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Table matérialisée (Créer un Lakehouse unifié pour le streaming et le traitement par lots)

Dernière mise à jour :Aug 09, 2026

Cette rubrique vous guide dans la création d'un pipeline analytique lakehouse unifié pour le streaming et le traitement par lots à l'aide de tables matérialisées. Vous apprendrez également à basculer du traitement par lots au traitement en continu en ajustant la fraîcheur des données d'une table matérialisée afin d'activer les mises à jour en temps réel.

Qu'est-ce qu'une table matérialisée ?

La table matérialisée est un nouveau type de table introduit dans Flink SQL pour simplifier les pipelines de données, tant pour le traitement par lots que pour le streaming, et offrir une expérience de développement unifiée. Lors de la création d'une table matérialisée, il n'est pas nécessaire de déclarer les champs ni leurs types. Spécifiez simplement la fraîcheur des données souhaitée et une instruction de requête. Le moteur Flink déduit automatiquement le schéma à partir de la requête et crée un pipeline d'actualisation des données correspondant pour respecter l'exigence de fraîcheur spécifiée. Pour plus d'informations, consultez Gestion des tables matérialisées.

Schéma du pipeline Lakehouse en temps réel

  1. Flink écrit les données provenant des sources dans Paimon pour former la couche ODS (Operational Data Store).

  2. Flink enrichit et élargit les données de la couche ODS via des jointures de tables, puis écrit le résultat dans une table matérialisée pour constituer la couche DWD (Data Warehouse Detail).

  3. Plusieurs tables matérialisées avec différents paramètres de fraîcheur effectuent des agrégations métier multidimensionnelles pour former la couche DWS (Data Warehouse Service), qui sert les requêtes applicatives.

Prérequis

  • Vous avez créé un espace de travail Flink. Pour plus d'informations, consultez Activer Realtime Compute for Apache Flink.

  • Si vous accédez aux ressources en tant qu'utilisateur ou rôle RAM (Resource Access Management), vérifiez que vous disposez des autorisations requises pour la console Flink. Pour plus d'informations, consultez Gestion des autorisations.

Étape 1 : Préparer les données de test

  1. Créez un catalogue Paimon.

    Les tables matérialisées sont alimentées par Apache Paimon. Vous devez créer un catalogue Paimon avec un type de métastore défini sur Filesystem . Ignorez cette étape si vous en possédez déjà un. Pour plus d'informations, consultez Créer un catalogue Paimon.

    Créer un catalogue Paimon

    1. Connectez-vous à la console de gestion Realtime Compute.

    2. Cliquez sur Console dans la colonne Actions de votre espace de travail cible.

    3. Dans le volet de navigation de gauche, sélectionnez Catalogs puis cliquez sur Create Catalog. Sélectionnez Apache Paimon et cliquez sur Next.

      Description des paramètres :

      Élément de configuration

      Description

      Remarques

      metastore

      Type de métastore.

      Cet exemple utilise filesystem comme type de métastore.

      catalog name

      Nom du catalogue Paimon.

      Saisissez un nom personnalisé en anglais. Cet exemple utilise paimon.

      warehouse

      Répertoire de l'entrepôt de données dans OSS.

      Utilisez le format oss://<bucket>/<object>. Où :

      • <bucket> : le nom de votre compartiment OSS.

      • <object> : le chemin d'accès où vos données sont stockées.

      Vérifiez les noms de votre compartiment et de votre objet dans la console de gestion OSS.

      fs.oss.endpoint

      Adresse de connexion à OSS.

      Si Flink et OSS se trouvent dans la même région, utilisez le point de terminaison du réseau privé. Sinon, utilisez le point de terminaison public. Pour plus d'informations, consultez Régions et points de terminaison.

      fs.oss.accessKeyId

      AccessKey ID de votre compte Alibaba Cloud ou utilisateur RAM disposant des autorisations de lecture et d'écriture sur OSS.

      Pour savoir comment l'obtenir, consultez Créer une clé d'accès. Pour éviter d'exposer les informations d'identification en clair, utilisez plutôt des variables. Pour plus d'informations, consultez Gestion des variables.

      fs.oss.accessKeySecret

      AccessKey Secret de votre compte Alibaba Cloud ou utilisateur RAM disposant des autorisations de lecture et d'écriture sur OSS.

  2. Créez la table des journaux de comportement utilisateur ods_user_log et la table d'informations produit ods_dim_product.

    1. Connectez-vous à la console de gestion Realtime Compute.

    2. Cliquez sur Console dans la colonne Actions de votre espace de travail cible.

    3. Dans le volet de navigation de gauche, sélectionnez Development > Scripts. Copiez et collez le code suivant pour créer les tables source.

      Cet exemple suppose que vous avez déjà créé un catalogue Paimon nommé paimon et que vous utilisez la base de données par défaut.
      CREATE TABLE `paimon`.`default`.`ods_user_log` (
        item_id INT NOT NULL,
        user_id INT NOT NULL,
        vtime TIMESTAMP(6),
        ds VARCHAR(10) NOT NULL
      ) 
      PARTITIONED BY(ds)
      WITH (
        'bucket' = '4',            -- Set the number of buckets to 4.
        'bucket-key' = 'item_id'   -- Specify the key used to determine bucket assignment. Rows with the same item_id go into the same bucket.
      );
      CREATE TABLE `paimon`.`default`.`ods_dim_product` (
        item_id INT NOT NULL,
        title VARCHAR(255),
        pict_url VARCHAR(255), 
        brand_id INT,
        seller_id INT,
        PRIMARY KEY(item_id) NOT ENFORCED
      ) WITH (
        'bucket' = '4',
        'bucket-key' = 'item_id'
      );
    4. Cliquez sur Run dans le coin supérieur droit pour créer les tables.

    5. Dans le volet de navigation de gauche, sélectionnez Catalogs, cliquez sur votre catalogue Paimon, puis cliquez sur Refresh pour afficher les nouvelles tables.

  3. Utilisez le connecteur Faker pour la génération de données simulées afin de générer des données fictives et de les écrire dans les tables Paimon.

    1. Dans le volet de navigation de gauche, sélectionnez Development > ETL.

    2. Cliquez sur New, sélectionnez Blank stream draft, cliquez sur Next, puis cliquez sur Create.

    3. Copiez l'instruction SQL suivante dans l'éditeur SQL.

      CREATE TEMPORARY TABLE `user_log` (
        item_id INT,  // Product ID
        user_id INT,  // User ID
        vtime TIMESTAMP,  
        ds AS DATE_FORMAT(CURRENT_DATE,'yyyyMMdd')
      ) WITH (
        'connector' = 'faker',    -- Faker connector
        'fields.item_id.expression'='#{number.numberBetween ''0'',''1000''}',    -- Generate a random number between 0 and 1000
        'fields.user_id.expression'='#{number.numberBetween ''0'',''100''}',
        'fields.vtime.expression'='#{date.past ''5'',''HOURS''}',           -- Generate data up to 5 hours before the current time
        'rows-per-second' = '3'   -- Generate 3 rows per second
       );
       CREATE TEMPORARY TABLE `dim_product` (
        item_id INT NOT NULL,
        title VARCHAR(255),
        pict_url VARCHAR(255), 
        brand_id INT,
        seller_id INT,
        PRIMARY KEY(item_id) NOT ENFORCED
       ) WITH (
        'connector' = 'faker',    -- Faker connector
        'fields.item_id.expression'='#{number.numberBetween ''0'',''1000''}',
        'fields.title.expression'='#{book.title}',
        'fields.pict_url.expression'='#{internet.domainName}',
        'fields.brand_id.expression'='#{number.numberBetween ''1000'',''10000''}',   
        'fields.seller_id.expression'='#{number.numberBetween ''1000'',''10000''}',
        'rows-per-second' = '3'        -- Generate 3 rows per second
       );
      BEGIN STATEMENT SET; 
      INSERT INTO `paimon`.`default`.`ods_user_log` 
        SELECT 
        item_id,
        user_id,
        vtime,
        CAST(ds AS VARCHAR(10)) AS ds
      FROM `user_log`;
      INSERT INTO `paimon`.`default`.`ods_dim_product`
        SELECT 
        item_id,
        title,
        pict_url,
        brand_id,
        seller_id
      FROM `dim_product`;
      END; 
    4. Cliquez sur Deploy dans le coin supérieur droit pour déployer la tâche.

    5. Dans le volet de navigation de gauche, sélectionnez O&M > Deployments. Cliquez sur Start dans la colonne Actions de votre tâche cible, sélectionnez Initial Mode, puis cliquez sur Start.

  4. Interrogez les données simulées.

    Dans le volet de navigation de gauche, sélectionnez Development > Scripts. Copiez l'instruction SQL suivante dans l'éditeur SQL et cliquez sur Run dans le coin supérieur droit.

    SELECT * FROM `paimon`.`default`.ods_dim_product LIMIT 10;
    SELECT * FROM `paimon`.`default`.ods_user_log LIMIT 10;

Étape 2 : Créer des tables matérialisées

Cette section permet de construire une table matérialisée de couche DWD nommée dwd_user_log_product en élargissant les tables source. Elle crée ensuite des tables matérialisées en aval basées sur dwd_user_log_product pour l'agrégation métier, complétant ainsi la couche DWS.

  1. Construisez la couche DWD de l'entrepôt de données en créant la table matérialisée dwd_user_log_product.

    1. Dans le volet de navigation de gauche, sélectionnez Catalogs et cliquez sur votre catalogue Paimon cible.

    2. Cliquez sur votre base de données cible (default dans cet exemple), puis cliquez sur Create Materialized Table. Copiez l'instruction SQL suivante dans l'éditeur SQL et cliquez sur Create.

      -- DWD layer widening logic
      CREATE MATERIALIZED TABLE dwd_user_log_product(
          PRIMARY KEY (item_id) NOT ENFORCED
      )
      PARTITIONED BY(ds)
      WITH (
        'partition.fields.ds.date-formatter' = 'yyyyMMdd'
      )
      FRESHNESS = INTERVAL '1' HOUR      -- Refresh every hour
      AS SELECT
        l.ds,
        l.item_id,
        l.user_id,
        l.vtime,
        r.brand_id,
        r.seller_id
      FROM `paimon`.`default`.`ods_user_log` l INNER JOIN `paimon`.`default`.`ods_dim_product` r
      ON l.item_id = r.item_id;
  2. Construisez la couche DWS en effectuant des agrégations métier multidimensionnelles basées sur la table matérialisée dwd_user_log_product.

    Cette rubrique explique comment créer la table matérialisée dws_overall qui agrège les comptes horaires de PV/UV par jour. Suivez l'étape précédente pour créer la table matérialisée dws_overall.

    // Aggregate PV/UV by day
    CREATE MATERIALIZED TABLE dws_overall(
        PRIMARY KEY(ds, hh) NOT ENFORCED
    )
    PARTITIONED BY(ds)
    WITH (
      'partition.fields.ds.date-formatter' = 'yyyyMMdd'
    )
    FRESHNESS = INTERVAL '1' HOUR   -- Refresh every hour
    AS SELECT 
        ds,
        COALESCE(hh, 'day') AS hh,
        count(*) AS pv,
        count(distinct user_id) AS uv
        FROM (SELECT ds, date_format(vtime, 'HH') AS hh, user_id 
    FROM `paimon`.`default`.`dwd_user_log_product`) tmp
    GROUP BY GROUPING SETS(ds, (ds, hh));

Étape 3 : Mettre à jour les tables matérialisées

Démarrer la mise à jour

La fraîcheur des données dans cet exemple est définie sur 1 heure. Après avoir cliqué sur Start, les mises à jour des données auront un retard d'au moins 1 heure par rapport aux mises à jour des tables de base.

  1. Dans le volet de navigation de gauche, sélectionnez O&M > Data Lineage et recherchez votre table matérialisée cible.

  2. Cliquez sur la vue de la table matérialisée, puis cliquez sur Start dans le coin inférieur droit de la page.

Remplissage des données

Le remplissage des données réécrit les données historiques dans des partitions spécifiques ou dans l'intégralité de la table. Utilisez cette fonction pour corriger les résultats du traitement en continu ou pour mettre à jour immédiatement les données destinées aux tâches par lots qui n'ont pas encore atteint leur heure planifiée.

Sélectionnez la vue de la table matérialisée dwd_user_log_product et cliquez sur Trigger Update dans le coin inférieur droit. Saisissez la date actuelle (par exemple, 20241216) comme nom de partition, cochez Cascade update downstream materialized table, puis cliquez sur Confirm. Dans la boîte de dialogue de confirmation, cliquez sur Confirm pour écraser les données immédiatement.

Pour plus d'informations sur le remplissage des données, consultez Remplir les données historiques.

Modifier la fraîcheur des données

Vous pouvez ajuster la fraîcheur des données pour mettre à jour les tables matérialisées quotidiennement, horairement, minutièrement, voire à la seconde, selon les besoins de votre activité.

Mettez à jour les paramètres de fraîcheur pour les tables matérialisées dwd_user_log_product et dws_overall. Cliquez sur la vue de la table matérialisée, puis cliquez sur Modify Freshness dans le coin inférieur droit. Définissez la fraîcheur au niveau de la minute pour des mises à jour en temps réel.

Pour plus d'informations sur la modification de la fraîcheur des données, consultez Modifier la fraîcheur des données.

Étape 4 : Interroger les tables matérialisées

Aperçu des données

Vous pouvez prévisualiser les 100 dernières lignes d'une table matérialisée.

  1. Dans le volet de navigation de gauche, sélectionnez O&M > Data Lineage et recherchez votre table matérialisée cible.

  2. Cliquez sur la vue de la table matérialisée, puis cliquez sur Details dans le coin inférieur droit de la page.

  3. Dans l'onglet Data Preview de la table matérialisée, cliquez sur l'icône Query.

Requête de données

Dans le volet de navigation de gauche, sélectionnez Development > Scripts. Copiez l'instruction SQL suivante dans l'éditeur SQL, sélectionnez l'extrait de code, puis cliquez sur Run pour interroger la table matérialisée dws_overall.

SELECT * FROM `paimon`.`default`.dws_overall ORDER BY hh;

Références