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
|
|
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 |
|
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 |
|
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.
|
Procédure
Connectez-vous à la console Realtime Compute for Apache Flink.
Dans la colonne Actions de l'espace de travail cible, cliquez sur Console.
Dans le volet de navigation de gauche, cliquez sur Catalogs, puis sur le catalogue Apache Paimon cible.
-
Cliquez sur la base de données cible, puis sur Create Materialized Table.
Supposons que vous disposiez d'une table source nommée
ordersavec une clé primaireorder_id, un nom de catégorieorder_nameet un champ de dateds. 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_orderbasée sur la tableorders. 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_idbasée sur la table matérialiséemt_order. La requête sélectionneorder_idetdscomme colonnes de table, définitorder_idcomme clé primaire,dscomme 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_dsbasée sur la table matérialiséemt_orderet spécifiez undate-formatter(format temporel) pour la colonne de partitionnementds. À chaque exécution d'une planification, l'heure de planification moins la fraîcheur est convertie en la valeur de partitiondscorrespondante. Par exemple, si la fraîcheur des données est définie sur 1 heure et que l'heure de planification est2024-01-01 00:00:00, la valeur calculée de ds est 20231231, et seules les données de la partitionds = '20231231'sont actualisées. Si l'heure planifiée est2024-01-01 01:00:00, la valeur calculée de ds est 20240101, et les données de la partitionds = '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 ;RemarqueDans
partition.fields.#.date-formatter, l'espace réservé#doit être une colonne de partitionnement valide de type STRING.L'option
partition.fields.#.date-formatterspé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.
-
-
Démarrez ou arrêtez la mise à jour de la table matérialisée.
Cliquez sur la table matérialisée cible sous son catalogue.
-
Dans le coin supérieur droit, cliquez sur Start ou Stop.
RemarqueSi vous arrêtez une mise à jour en cours, le job termine le cycle de mise à jour actuel avant de s'arrêter.
-
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
-
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 ; -
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.quantityAStotal_price, et inclut une nouvelle condition WHEREorders.price*orders.quantity> 1000. Dans la zone Table Columns, le nouveau champtotal_price(DOUBLE) est mis en surbrillance en vert. Cliquez sur OK. Vous pouvez afficher la colonne nouvellement ajoutée et la logique de requête modifiée sous l'onglet Table Schema.
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-formattern'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>etWHERE <column> [NOT] IN <subquery>ne sont pas prises en charge.UNION
Seul
UNION ALLest pris en charge.JOIN
-
INNER JOINest pris en charge. -
LEFT/RIGHT/FULL [OUTER] JOINne sont pas pris en charge, sauf dans les cas spécifiques deLATERAL JOINet de lookup join décrits ci-dessous. -
[LEFT [OUTER]] JOIN LATERALavec 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 queSELECT * FROM a, b WHERE a.id = b.id, sont prises en charge. -
Le calcul incrémentiel pour
INNER JOINlit 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.
Cliquez sur la table matérialisée cible sous son catalogue.
-
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 partitionds=20241201font 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.
RemarqueLors 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.
RemarqueLors 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.
-
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
Sous le catalogue correspondant, cliquez sur la base de données materialized table, puis sur la materialized table cible.
-
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 à 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
Pour une introduction aux tables matérialisées, consultez la section Gestion des tables matérialisées.
Pour savoir comment construire un pipeline analytique dans un data lakehouse avec un traitement batch et streaming unifié basé sur Paimon et les tables matérialisées, et comment passer du mode batch au mode streaming en modifiant la fraîcheur des données pour des mises à jour en temps réel, consultez la section Démarrage rapide : construction d'un data lakehouse avec un traitement batch et streaming unifié à l'aide de tables matérialisées.