Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Création et utilisation de tables matérialisées

Dernière mise à jour :Aug 09, 2026

Cette rubrique explique comment créer une table matérialisée, rétrocharger les données, modifier la fraîcheur des données et consulter la lignée des données.

Limites

  • Cette fonctionnalité est disponible uniquement dans Ververica Runtime (VVR) 8.0.10 et versions ultérieures.

  • Les tables matérialisées peuvent être créées uniquement dans un catalogue Apache Paimon qui utilise Filesystem ou DLF pour le stockage des métadonnées. Les catalogues Apache Paimon personnalisés ne sont pas pris en charge.

  • Vous devez disposer des autorisations nécessaires pour développer et déployer des jobs. Pour plus d'informations, consultez la section Autorisation de la console de développement.

  • Les objets temporaires, tels que les tables temporaires, les fonctions définies par l'utilisateur temporaires et les vues temporaires, ne sont pas pris en charge.

Création d'une table matérialisée

Syntaxe

CREATE MATERIALIZED TABLE [catalog_name.][db_name.]table_name
-- Primary key constraint
[([CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED)]
[COMMENT table_comment]
-- Partition key
[PARTITIONED BY (partition_column_name1, partition_column_name2, ...)]
-- With options
[WITH (key1=val1, key2=val2, ...)]
-- Data freshness
FRESHNESS = INTERVAL '<num>' { SECOND | MINUTE | HOUR | DAY }
-- Refresh mode
[REFRESH_MODE = { CONTINUOUS | FULL }]
AS  <select_statement>

Paramètres

Paramètre

Obligatoire

Description

FRESHNESS

Oui

Fraîcheur des données de la table matérialisée, définissant la latence maximale autorisée pour les mises à jour des données provenant des tables source.

Remarque
  • Si une table en amont est également une table matérialisée, la fraîcheur des données de la table en aval doit être un multiple entier positif de la fraîcheur des données de la table en amont.

  • La fraîcheur maximale des données est d'un jour.

AS <select_statement>

Oui

Définit la requête qui alimente la table matérialisée. La table en amont peut être une table matérialisée, une table standard ou une vue. L'instruction SELECT prend en charge toutes les requêtes Flink SQL.

PRIMARY KEY

Non

Colonnes facultatives identifiant de manière unique chaque ligne de la table. Ces colonnes ne peuvent pas contenir de valeurs nulles.

PARTITIONED BY

Non

Colonnes facultatives utilisées pour partitionner la table matérialisée.

Options WITH

Non

Définit les propriétés de la table et les paramètres de format temporel pour les colonnes de partitionnement.

Par exemple, vous pouvez définir le paramètre de format temporel pour une colonne de partitionnement avec WITH ('partition.fields.#.date-formatter' = 'yyyyMMdd'). Pour plus d'informations sur l'utilisation des paramètres, consultez les exemples de la procédure.

REFRESH_MODE

Non

Spécifie le mode d'actualisation de la table matérialisée. Un mode d'actualisation spécifié est prioritaire sur le mode inféré automatiquement par le framework en fonction de la fraîcheur des données. Cela permet de gérer des scénarios spécifiques.

  • CONTINUOUS : un job de streaming met à jour la table matérialisée de manière incrémentielle. Les données en aval deviennent visibles immédiatement ou après la fin d'un point de contrôle.

  • FULL : un workflow déclenche périodiquement des mises à jour par lots pour la table matérialisée. Dans ce mode, le moteur détermine automatiquement s'il convient d'effectuer une mise à jour complète ou incrémentielle. Pour plus d'informations, consultez la section Mises à jour incrémentielles d'une table matérialisée. Le cycle d'actualisation des données correspond au paramètre de fraîcheur des données. Par défaut, les données sont écrasées au niveau de la table. Si des colonnes de partitionnement existent, vous pouvez choisir d'actualiser uniquement la dernière partition ou de mettre à jour toutes les partitions.

Procédure

  1. Connectez-vous à la console Realtime Compute for Apache Flink.

  2. Dans la colonne Actions de l'espace de travail cible, cliquez sur Console.

  3. Dans le volet de navigation de gauche, cliquez sur Catalogs, puis sur le catalogue Apache Paimon cible.

  4. Cliquez sur la base de données cible, puis sur Create Materialized Table.

    Supposons que vous disposiez d'une table source nommée orders avec une clé primaire order_id, un nom de catégorie order_name et un champ de date ds. Les exemples suivants montrent comment créer des tables matérialisées basées sur cette table :

    • Créez une table matérialisée mt_order basée sur la table orders. La requête sélectionne toutes les colonnes et la fraîcheur des données est définie sur 5 secondes.

      CREATE MATERIALIZED TABLE mt_order
      FRESHNESS = INTERVAL '5' SECOND
      AS
      SELECT * FROM `paimon`.`db`.`orders`
      ;
    • Créez une table matérialisée mt_id basée sur la table matérialisée mt_order. La requête sélectionne order_id et ds comme colonnes de table, définit order_id comme clé primaire, ds comme colonne de partitionnement et la fraîcheur des données sur 30 minutes.

      CREATE MATERIALIZED TABLE mt_id (
       PRIMARY KEY (order_id) NOT ENFORCED
      )
      PARTITIONED BY(ds)
      FRESHNESS = INTERVAL '30' MINUTE
      AS
      SELECT order_id,ds FROM mt_order
      ;
    • Créez la table matérialisée mt_ds basée sur la table matérialisée mt_order et spécifiez un date-formatter (format temporel) pour la colonne de partitionnement ds. À chaque exécution d'une planification, l'heure de planification moins la fraîcheur est convertie en la valeur de partition ds correspondante. Par exemple, si la fraîcheur des données est définie sur 1 heure et que l'heure de planification est 2024-01-01 00:00:00, la valeur calculée de ds est 20231231, et seules les données de la partition ds = '20231231' sont actualisées. Si l'heure planifiée est 2024-01-01 01:00:00, la valeur calculée de ds est 20240101, et les données de la partition ds = '20240101' sont actualisées.

      CREATE MATERIALIZED TABLE mt_ds
      PARTITIONED BY(ds)
      WITH (
          'partition.fields.ds.date-formatter' = 'yyyyMMdd'
      )
      FRESHNESS = INTERVAL '1' HOUR
      AS
      SELECT order_id,order_name,ds FROM mt_order
      ;
      Remarque
      • Dans partition.fields.#.date-formatter, l'espace réservé # doit être une colonne de partitionnement valide de type STRING.

      • L'option partition.fields.#.date-formatter spécifie le format de partition temporelle pour la table matérialisée. L'espace réservé # représente le nom d'une colonne de partitionnement de type chaîne. Cette information permet au système d'identifier la partition à actualiser lors d'une mise à jour planifiée.

  5. Démarrez ou arrêtez la mise à jour de la table matérialisée.

    1. Cliquez sur la table matérialisée cible sous son catalogue.

    2. Dans le coin supérieur droit, cliquez sur Start ou Stop.

      Remarque

      Si vous arrêtez une mise à jour en cours, le job termine le cycle de mise à jour actuel avant de s'arrêter.

  6. Consultez les détails du job de la table matérialisée.

    Sous l'onglet Table Schema, dans la section Basic Information, cliquez sur l'ID du job situé à côté de Latest Job ou de Workflow pour afficher les détails.

Modification d'une requête de table matérialisée

Limites

  • Vous pouvez modifier la requête uniquement pour les tables matérialisées créées dans VVR 11.1 ou version ultérieure.

  • Lors de la modification d'une requête, vous pouvez uniquement ajouter des colonnes et modifier la logique de calcul. Vous ne pouvez pas modifier l'ordre des colonnes existantes ni leurs définitions.

    Opération

    Prise en charge

    Description

    Ajouter une nouvelle colonne

    Oui

    Vous pouvez ajouter de nouvelles colonnes au schéma tout en conservant l'ordre des colonnes existantes.

    Modifier la logique de calcul d'une colonne existante (sans changer son nom ni son type)

    Oui

    Vous pouvez modifier la logique de calcul, mais le nom de la colonne et le type de données doivent rester identiques.

    Modifier l'ordre des colonnes existantes

    Non

    L'ordre des colonnes est fixe. Pour le modifier, vous devez supprimer et recréer la table matérialisée.

    Modifier le nom ou le type de données d'une colonne existante

    Non

    Vous devez supprimer et recréer la table matérialisée.

Exemple

  1. Cliquez sur Edit Table et modifiez la requête. Le code suivant fournit un exemple :

    ALTER MATERIALIZED TABLE `paimon`.`default`.`mt-orders`
        AS
        SELECT
          *,
          price * quantity AS total_price
        FROM orders
        WHERE price * quantity > 1000
    ;
  2. Cliquez sur Preview pour afficher une comparaison avant/après.

    La boîte de dialogue Materialized Table Change Details affiche les modifications suivantes : dans la zone SQL Statement, l'instruction SQL passe de la commande originale SELECT * à SELECT *, orders.price * orders.quantity AS total_price, et inclut une nouvelle condition WHERE orders.price * orders.quantity > 1000. Dans la zone Table Columns, le nouveau champ total_price (DOUBLE) est mis en surbrillance en vert.

  3. Cliquez sur OK. Vous pouvez afficher la colonne nouvellement ajoutée et la logique de requête modifiée sous l'onglet Table Schema.

Important

L'ajout de colonnes n'affecte généralement pas les consommateurs en aval. Toutefois, si un job en aval repose sur une analyse dynamique (telle que SELECT * ou le mappage automatique des champs) lors de la consommation des données de la table matérialisée en amont, le job risque d'échouer ou de signaler une erreur d'incompatibilité de format de données. Nous vous recommandons d'éviter l'analyse dynamique, d'utiliser des listes de colonnes fixes et de mettre à jour rapidement le schéma de la table en aval chaque fois que le schéma en amont change.

Mises à jour incrémentielles

Limites

Cette fonctionnalité est disponible uniquement dans VVR 8.0.11 et versions ultérieures.

Modes de mise à jour

Les tables matérialisées prennent en charge trois modes de mise à jour : streaming, lot complet et lot incrémentiel.

Le mode est déterminé par le paramètre de fraîcheur des données. Une fraîcheur inférieure à 30 minutes active le mode streaming, tandis qu'une fraîcheur de 30 minutes ou plus active le mode lot. En mode lot, le moteur décide automatiquement entre une mise à jour complète ou incrémentielle. Une mise à jour incrémentielle calcule uniquement les données qui ont changé depuis la dernière mise à jour et les fusionne dans la table matérialisée. Une mise à jour complète calcule les données pour l'ensemble de la table ou de la partition et écrase les données existantes dans la table matérialisée. En mode lot, le moteur privilégie les mises à jour incrémentielles et ne revient à une mise à jour complète que lorsqu'une mise à jour incrémentielle n'est pas possible.

Conditions de mise à jour incrémentielle

Une mise à jour incrémentielle est effectuée uniquement si la table matérialisée remplit toutes les conditions suivantes :

  • Le paramètre partition.fields.#.date-formatter n'est pas configuré dans la définition de la table.

  • La table source ne possède pas de clé primaire définie.

  • La requête dans la définition de la table matérialisée prend en charge les mises à jour incrémentielles comme décrit dans le tableau suivant :

    Clause SQL

    Prise en charge

    SELECT

    Prise en charge de la sélection de colonnes et des expressions de fonctions scalaires, y compris les fonctions définies par l'utilisateur. Les fonctions d'agrégation ne sont pas prises en charge.

    FROM

    Prise en charge des noms de tables et des sous-requêtes.

    WITH

    Prise en charge des Common Table Expressions (CTE).

    WHERE

    Prise en charge des conditions de filtre incluant diverses expressions de fonctions scalaires, y compris les fonctions définies par l'utilisateur. Les sous-requêtes telles que WHERE [NOT] EXISTS <subquery> et WHERE <column> [NOT] IN <subquery> ne sont pas prises en charge.

    UNION

    Seul UNION ALL est pris en charge.

    JOIN

    • INNER JOIN est pris en charge.

    • LEFT/RIGHT/FULL [OUTER] JOIN ne sont pas pris en charge, sauf dans les cas spécifiques de LATERAL JOIN et de lookup join décrits ci-dessous.

    • [LEFT [OUTER]] JOIN LATERAL avec une expression de fonction de table (y compris les fonctions définies par l'utilisateur) est pris en charge.

    • Pour le lookup join, seul A [LEFT [OUTER]] JOIN B FOR SYSTEM_TIME AS OF PROCTIME() est pris en charge.

    Remarque
    • Les jointures implicites sans le mot-clé JOIN, telles que SELECT * FROM a, b WHERE a.id = b.id, sont prises en charge.

    • Le calcul incrémentiel pour INNER JOIN lit toujours l'intégralité des données des deux tables source.

    GROUP BY

    Non pris en charge.

Exemples de mises à jour incrémentielles

Exemple 1 : Traitez les données de la table source orders à l'aide de fonctions scalaires.

CREATE MATERIALIZED TABLE mt_shipped_orders (
    PRIMARY KEY (order_id) NOT ENFORCED
)
FRESHNESS = INTERVAL '30' MINUTE
AS
SELECT 
    order_id,
    COALESCE(customer_id, 'Unknown') AS customer_id,
    CAST(order_amount AS DECIMAL(10, 2)) AS order_amount,
    CASE 
        WHEN status = 'shipped' THEN 'Completed'
        WHEN status = 'pending' THEN 'In Progress'
        ELSE 'Unknown'
    END AS order_status,
    DATE_FORMAT(order_ts, 'yyyyMMdd') AS order_date,
    UDSF_ProcessFunction(notes) AS notes
FROM 
    orders
WHERE
    status = 'shipped';

Exemple 2 : Enrichissez les données de la table source orders à l'aide de LATERAL JOIN et de lookup join.

CREATE MATERIALIZED TABLE mt_enriched_orders (
    PRIMARY KEY (order_id, order_tag) NOT ENFORCED
)
FRESHNESS = INTERVAL '30' MINUTE
AS
WITH o AS (
    SELECT
        order_id,
        product_id,
        quantity,
        proc_time,
        e.tag AS order_tag
    FROM 
        orders,
        LATERAL TABLE(UDTF_StringSplitFunction(tags, ',')) AS e(tag))
SELECT 
    o.order_id,
    o.product_id,
    p.product_name,
    p.category,
    o.quantity,
    p.price,
    o.quantity * p.price AS total_amount,
    order_tag
FROM o 
LEFT JOIN 
    product_info FOR SYSTEM_TIME AS OF PROCTIME() AS p
ON 
    o.product_id = p.product_id;

Rétrochargement des données

Auparavant, la correction des résultats du traitement de flux pour les données historiques nécessitait le développement d'un job batch distinct. En utilisant des tables matérialisées, vous pouvez effectuer directement un rétrochargement des partitions de données historiques. Cette approche unifie le traitement batch et streaming, ce qui réduit les coûts de développement et d'exploitation.

  1. Cliquez sur la table matérialisée cible sous son catalogue.

  2. Sous l'onglet Data Information, effectuez le rétrochargement des données.

    Si vous avez défini une colonne de partitionnement lorsque vous avez créé la table matérialisée, il s'agit d'une table partitionnée. Sinon, il s'agit d'une table non partitionnée.

    Table partitionnée

    Dans la section Partitions, cliquez sur Trigger Update s'il s'agit de la première fois que vous effectuez un rétrochargement des données ou si la partition requise n'existe pas. Si les partitions existent déjà, vous pouvez sélectionner une partition spécifique et cliquer sur Refresh dans la colonne Actions.

    Cliquez sur l'onglet Data Information et effectuez les opérations dans la zone Data Partitions en bas de la page.

    Paramètres

    • Colonne de partitionnement : la colonne de partitionnement de la table. Par exemple, si vous saisissez 20241201, toutes les données de la partition ds=20241201 font l'objet d'un rétrochargement.

    • Nom de la tâche : le nom de la tâche de rétrochargement des données.

    • Plage de mise à jour (facultatif) : indique s'il faut propager les mises à jour aux tables matérialisées en aval. À partir de la table actuelle, toutes les tables matérialisées de la lignée des données sont mises à jour. La profondeur maximale prise en charge de la lignée en aval est de 6 niveaux.

      Remarque
      • Lors de la mise à jour d'une table partitionnée, les tables matérialisées en aval doivent avoir exactement les mêmes colonnes de partitionnement que la table de départ. Sinon, l'opération de mise à jour échoue.

      • Si une mise à jour échoue pour une table matérialisée quelconque de la lignée, tous les nœuds en aval suivants échouent également.

    • Cible de déploiement : vous pouvez sélectionner une file d'attente ou un cluster de session. La sélection par défaut est default-queue.

    Table non partitionnée

    Dans la section Data Detail, cliquez sur Refresh.

    Paramètres

    • Nom de la tâche : le nom de la tâche de rétrochargement des données.

    • Plage de mise à jour : cette option n'est pas disponible pour les tables non partitionnées.

      Remarque
      • Lors d'une mise à jour, les données en aval sont entièrement actualisées.

      • Si une mise à jour échoue pour une table matérialisée quelconque de la lignée, tous les nœuds en aval suivants échouent également.

      • Les mises à jour en cascade ne sont pas prises en charge si la table de départ est une table non partitionnée mise à jour par un job de streaming.

    • Cible de déploiement : vous pouvez sélectionner une file d'attente ou un cluster de session. La sélection par défaut est default-queue.

  3. Planification et rétrochargement par lots.

    Vous pouvez utiliser Task Orchestration pour créer un workflow pour la table matérialisée afin d'exécuter des jobs de rétrochargement selon une planification. Vous pouvez également utiliser la fonctionnalité de rétrochargement des données du workflow pour effectuer un rétrochargement des données en masse pour une plage horaire spécifiée.

Modification de la fraîcheur des données

  1. Sous le catalogue correspondant, cliquez sur la base de données materialized table, puis sur la materialized table cible.

  2. Dans le coin supérieur droit, cliquez sur Modify Freshness.

    • Si la table matérialisée ne possède pas de clé primaire, vous ne pouvez pas basculer sa méthode de mise à jour entre streaming et lot. Le système utilise un job de streaming pour les valeurs de fraîcheur des données inférieures à 30 minutes et un job batch pour les valeurs de 30 minutes ou plus. Par conséquent, le franchissement de ce seuil de 30 minutes n'est pas autorisé pour les tables sans clé primaire.

    • Si la table en amont est une table matérialisée, assurez-vous que la fraîcheur des données de la table en aval est un multiple entier positif de la fraîcheur des données de la table en amont.

    • La fraîcheur maximale des données est d'un jour.

Consultation de la lignée des données

Dans le volet de navigation de gauche, accédez à O&M > Data Lineage pour accéder à la page de lignée des données pour les tables matérialisées. Sur cette page, vous pouvez afficher les relations de lignée entre toutes les tables matérialisées. Vous pouvez également effectuer des opérations telles que Start/Stop Update et Modify Freshness directement sur une table matérialisée. Cliquez sur Details pour accéder à la page de détails de la table matérialisée correspondante.

Cliquez sur un nœud de table matérialisée pour développer son panneau de détails. Le panneau affiche la valeur Data Freshness, l'heure de la Latest Update Time et le Materialized Table Status, et propose une action Trigger Update.

Documents connexes