Meilleures pratiques pour LakeFlow Pipelines
Appliquez ces modèles recommandés lors de la conception, de la création et de l'exploitation de pipelines, que vous démarriez un nouveau pipeline ou que vous en amélioriez un existant.
Les pages suivantes approfondissent des décisions de conception spécifiques :
Sujet | Description |
|---|---|
Modélisez les données de la couche Gold sous forme de faits et de dimensions dans un schéma en étoile, et mappez la modélisation dimensionnelle sur les types de datasets de pipeline. | |
Comprenez l'idempotence et les cas où vous bénéficiez d'un traitement « exactly-once » default, par rapport aux situations où vous devez ajouter des écritures idempotentes. | |
Décidez combien de jeux de données appartiennent à un seul pipeline, et quand diviser le travail en pipelines séparés. | |
Parcourez une liste de contrôle concernant la qualité, la fiabilité, l’observabilité, le déploiement, le coût et la gouvernance des données avant d’exécuter un pipeline sans surveillance. |
Choisir le type de dataset approprié
Les pipelines proposent trois types de dataset : les tables de streaming, les vues matérialisées et les vues temporaires. Choisir le type approprié pour chaque couche de votre pipeline permet d'éviter des coûts de compute inutiles et facilite la compréhension de votre code.
Les streaming tables sont le choix idéal pour l'ingestion de données et les Transformations en streaming à faible latence. Chaque ligne d'entrée est lue et traitée une seule fois, ce qui les rend idéales pour les charges de travail en ajout seul, les données à haut volume et le traitement piloté par les événements à partir du stockage cloud ou des bus de messages.
Les vues matérialisées sont le choix idéal pour les transformations complexes et les requêtes analytiques. Leurs résultats sont pré-calculés et maintenus à jour grâce à un refresh incrémental, ce qui permet de les interroger rapidement. Vous ne pouvez pas modifier directement les données dans une vue matérialisée ; la définition de query contrôle la sortie.
Les vues temporaires sont des vues limitées au pipeline qui organisent votre logique de Transformations sans matérialiser de données dans le stockage. Utilisez-les pour les étapes intermédiaires qui n’ont pas besoin de leur propre table.
Le tableau suivant résume quand utiliser chaque type :
Cas d'usage | Type recommandé | Motif |
|---|---|---|
Ingestion à partir d’un stockage cloud ou d’un bus de messages | Table de streaming | Traite chaque enregistrement une fois ; gère les charges de travail à fort volume et uniquement les charges. |
CDC Stream (insertions, mises à jour, suppressions) | Table de streaming | Utilisé comme cible de |
Agrégations et jointures complexes | Vue matérialisée | Rafraîchie progressivement ; Évite le recalcul complet à chaque mise à jour. |
Accélération des query de tableau de bord | Vue matérialisée | Les résultats précalculés rendent les requêtes plus rapides que celles effectuées sur des tables brutes. |
Transformations intermédiaires (aucun lecteur en aval) | Vue temporaire | Organise la logique du pipeline sans encourir de coût de stockage. |
Pour plus d'informations, consultez les tables de streaming, les vues matérialisées et Qu'est-ce que les LakeFlow Pipelines ?.
Utilisez la CDC déclarative au lieu de MERGE impératif
L’implémentation de la capture des données de changement (CDC) avec des instructions SQL MERGE impératives nécessite un code personnalisé important pour gérer correctement l’ordre des événements, la déduplication, les mises à jour partielles et l’évolution des schémas. Chacune de ces préoccupations doit être résolue indépendamment, et le code qui en résulte est difficile à maintenir et à tester.
Les pipelines fournissent l’instruction AUTO CDC ... INTO (SQL) et la fonction create_auto_cdc_flow() (Python), qui gèrent de manière déclarative le classement, la déduplication, les événements en désordre et l’évolution des schémas. Vous décrivez la forme du flux de données de modification et de la table cible, et le pipeline s’occupe du reste. AUTO CDC prend en charge à la fois le SCD de type 1 (remplacement) et le SCD de type 2 (conservation de l’historique).
Pour plus d'informations, consultez Change data capture et instantanés et Les API AUTO CDC : simplifiez la capture des modifications de données avec les pipelines.
Assurez la qualité des données avec des attentes
Les attentes sont des expressions SQL vrai/faux appliquées à chaque ligne passant dans un dataset. Lorsqu’une ligne échoue à la condition, le pipeline répond selon la politique de violation que vous avez configurée. Les attentes envoient des métriques au journal des événements du pipeline, quelle que soit la politique, ce qui vous permet de suivre les tendances de la qualité des données au fil du temps.
Choisir une politique de violation
Trois politiques de violation sont disponibles. Choisissez celui qui correspond à votre tolérance aux données de mauvaise qualité :
- avertissement (default) : Les enregistrements non valides sont écrits dans la table cible et signalés dans des métriques. Utilisez cette politique lorsque vous devez collecter toutes les données mais souhaitez avoir une visibilité sur les questions de qualité.
- abandonner : Les documents non valides sont jetés avant d’être écrits. Utilisez-le lorsque des mauvaises lignes sont attendues et ne devraient pas se propager en aval.
- fail : la mise à jour du pipeline s’arrête au premier enregistrement non valide. Utilisez cette option pour les données critiques où tout enregistrement incorrect indique un problème sérieux en amont.
Les exemples suivants montrent chaque politique appliquée à une table de streaming :
- SQL
- Python
-- Warn: write invalid records but track them in metrics
CREATE OR REFRESH STREAMING TABLE orders_raw (
CONSTRAINT valid_order_id EXPECT (order_id IS NOT NULL)
) AS SELECT * FROM STREAM read_files("/volumes/raw/orders", format => "json");
-- Drop: discard invalid records before writing
CREATE OR REFRESH STREAMING TABLE orders_clean (
CONSTRAINT non_negative_amount EXPECT (amount >= 0) ON VIOLATION DROP ROW
) AS SELECT * FROM STREAM(orders_raw);
-- Fail: stop the pipeline on any invalid record
CREATE OR REFRESH STREAMING TABLE orders_critical (
CONSTRAINT required_customer_id EXPECT (customer_id IS NOT NULL) ON VIOLATION FAIL UPDATE
) AS SELECT * FROM STREAM(orders_clean);
from pyspark import pipelines as dp
# Warn: write invalid records but track them in metrics
@dp.table
@dp.expect("valid_order_id", "order_id IS NOT NULL")
def orders_raw():
return spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.load("/volumes/raw/orders")
# Drop: discard invalid records before writing
@dp.table
@dp.expect_or_drop("non_negative_amount", "amount >= 0")
def orders_clean():
return spark.readStream.table("orders_raw")
# Fail: stop the pipeline on any invalid record
@dp.table
@dp.expect_or_fail("required_customer_id", "customer_id IS NOT NULL")
def orders_critical():
return spark.readStream.table("orders_clean")
Mettre en quarantaine les enregistrements non valides
Lorsque vous souhaitez conserver les enregistrements supprimés pour enquête plutôt que de les ignorer silencieusement, utilisez un modèle de quarantaine. Acheminez les lignes qui échouent à la validation vers une table de streaming distincte en utilisant deux flux : l’un qui supprime les lignes non valides de la table principale et un second qui écrit uniquement les lignes non valides dans une table de quarantaine. Cela vous permet d’examiner, de corriger et de retraiter les données erronées sans contaminer votre dataset propre.
Pour un exemple détaillé du modèle de quarantaine, consultez Recommandations d'attentes et modèles avancés.
Pour plus d'informations sur les attentes, consultez Gérer la qualité des données avec les attentes de pipeline.
Paramétrer vos pipelines
Les pipelines ont des paramètres default de catalogue et de schéma, donc le code qui lit et écrit dans le même catalogue et schéma fonctionne entre environnements sans aucun paramètre. Cependant, si votre pipeline doit référencer un second catalogue ou schéma (par exemple, lire un catalogue source partagé qui diffère entre développement et production), évitez de coder ces noms directement dans votre code source. À la place, définissez-les comme des paramètres de configuration du pipeline (paires clé-valeur définies dans les paramètres du pipeline) et référez-les dans votre code. Cela permet à une seule base de code de fonctionner correctement dans différents environnements en inversant les valeurs des paramètres.
- SQL
- Python
CREATE OR REFRESH MATERIALIZED VIEW transaction_summary AS
SELECT account_id, COUNT(txn_id) AS txn_count, SUM(amount) AS total_amount
FROM ${source_catalog}.sales.transactions
GROUP BY account_id;
from pyspark import pipelines as dp
from pyspark.sql.functions import count, sum
@dp.materialized_view
def transaction_summary():
source_catalog = spark.conf.get("source_catalog")
return spark.read.table(f"{source_catalog}.sales.transactions") \
.groupBy("account_id") \
.agg(
count("txn_id").alias("txn_count"),
sum("amount").alias("total_amount")
)
Pour plus d’informations, voir Utiliser des paramètres avec des pipelines.
Choisissez entre le mode Trigger et le mode pipeline continu
Le mode déclenché traite toutes les données disponibles puis s’arrête. C'est le bon choix pour la grande majorité des pipelines : ceux qui s'exécutent selon une planification (horaire, quotidienne ou à la demande) et qui ne nécessitent pas une fraîcheur des données inférieure à la minute.
Le mode continu maintient le cluster en cours d'exécution et traite les nouvelles données dès qu'elles arrivent. Cela n'est approprié que si votre cas d'usage nécessite une latence de l'ordre de quelques secondes à quelques minutes. Comme le mode continu nécessite un cluster toujours activé, il est nettement plus coûteux que le mode Trigger.
Le mode temps réel s'appuie sur le mode continu pour atteindre une latence inférieure à la seconde, de l'ordre de la milliseconde, pour les charges de travail opérationnelles telles que la détection de la fraude ou la personnalisation en temps réel. Cela nécessite une configuration et une planification du compute supplémentaires. Voir Utiliser le mode temps réel dans les LakeFlow Pipelines.
Pour plus d'informations, consultez Mode de pipeline Trigger ou continu et Configurer des pipelines.
Utilisez le clustering fluide pour le layout des données
Le clustering liquide remplace le partitionnement statique et ZORDER pour optimiser la Layout des données dans les tables Delta. Le partitionnement statique nécessite de sélectionner des colonnes de partition et de réorganiser les données au préalable, ce qui peut entraîner une asymétrie des données pour les valeurs inégalement réparties. Le clustering liquide est auto-ajustable, résistant à l'asymétrie et incrémentiel, ne réécrivant que les données nécessitant une réorganisation à chaque exécution.
Modifiez les colonnes de clustering à tout moment sans réécrire la table complète à mesure que les modèles de query évoluent.
Databricks recommande le clustering liquide automatique, qui permet à Databricks de sélectionner et de maintenir les colonnes de clustering optimales en fonction de votre charge de travail de query. Activez-le avec CLUSTER BY AUTO:
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE events
CLUSTER BY AUTO
AS SELECT * FROM STREAM read_files("/volumes/raw/events", format => "parquet");
from pyspark import pipelines as dp
@dp.table(cluster_by_auto=True)
def events():
return spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "parquet") \
.load("/volumes/raw/events")
Pour choisir vous-même les colonnes de regroupement, spécifiez-les explicitement :
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE events
CLUSTER BY (event_date, region)
AS SELECT * FROM STREAM read_files("/volumes/raw/events", format => "parquet");
from pyspark import pipelines as dp
@dp.table(cluster_by=["event_date", "region"])
def events():
return spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "parquet") \
.load("/volumes/raw/events")
Pour plus d’informations, voir Tables en streaming et Utiliser le cluster liquide pour les tableaux.
Gérer les pipelines avec la CI/CD et les Declarative Automation Bundles
Contrôlez les versions du code source de votre pipeline et utilisez des Déclarative Automation Bundles pour gérer les déploiements entre environnements.
Pour plus d’informations, consultez Créer un pipeline contrôlé par les sources, Convertir un pipeline en projet de bundle et Utiliser des paramètres avec des pipelines.
Stockez le code du pipeline dans un contrôle de version
Stockez tous les fichiers source du pipeline (Python et SQL) avec la configuration de votre bundle dans un repository Git. Le contrôle de version de l'intégralité du projet vous offre un historique complet des modifications, facilite la collaboration et vous permet de valider les changements dans un environnement de développement avant de les promouvoir en production.
Databricks recommande les Declarative Automation Bundles pour gérer ce flux de travail. Un bundle définit la configuration de votre pipeline en YAML aux côtés de votre code source, et la databricks bundle CLI vous permet de valider, déployer et exécuter des pipelines depuis votre terminal ou un système CI/CD.
Utilisez des cibles de bundle pour l'isolation de l'environnement
Les bundles permettent de définir plusieurs cibles (par exemple, dev, staging, prod), chacune avec son propre ensemble de remplacements pour les noms de catalogue, les politiques de cluster, les adresses de notification et d'autres paramètres. Combinez les cibles de bundle avec les parameter de pipeline pour injecter les valeurs correctes spécifiques à l'environnement au moment du déploiement, afin de garder votre code source exempt de constantes d'environnement.
Un workflow classique se présente comme suit :
- Un développeur travaille sur une Branch fonctionnalité, en déployant dans un pipeline de développement personnel dans un catalogue de développement.
- On Merge to the main Branch, a CI system runs
databricks bundle validateanddatabricks bundle deploy --target stagingto validate and deploy the pipeline to a staging environment. - Une fois les tests réussis, le système CI déploie en production avec
databricks bundle deploy --target prod.
Bonnes pratiques en streaming
Utilisez ces modèles pour gérer l’état, contrôler les données en retard et assurer la fiabilité des pipelines de streaming.
Pour plus d’informations, consultez Optimiser le traitement avec état avec des watermarks, Récupérer un pipeline après un échec de point de contrôle de streaming et Remplir les données historiques avec des pipelines.
Utiliser des filigranes pour les opérations avec état
Les filigranes (watermarks) limitent l'état que le pipeline conserve en mémoire lors des opérations de streaming avec état, telles que les agrégations par fenêtre et la déduplication. Sans filigrane, l'état augmente de manière illimitée à mesure que le pipeline accumule des données pour chaque clé possible, provoquant à terme des erreurs de mémoire insuffisante sur les pipelines à exécution longue.
Un filigrane (watermark) spécifie une colonne de Timestamp et un threshold de tolérance pour les données arrivant en retard. Les enregistrements qui arrivent après le dépassement du threshold sont ignorés. Choisissez un threshold qui équilibre votre tolérance aux données en retard par rapport au coût en mémoire lié au maintien de cet état ouvert.
L'exemple suivant calcule une agrégation par fenêtre basculante d'une minute avec un watermark de trois minutes :
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE event_counts AS
SELECT window(event_time, '1 minute') AS time_window, region, COUNT(*) AS cnt
FROM STREAM(events_raw)
WATERMARK event_time DELAY OF INTERVAL 3 MINUTES
GROUP BY time_window, region;
from pyspark import pipelines as dp
from pyspark.sql.functions import window
@dp.table
def event_counts():
return (
spark.readStream.table("events_raw")
.withWatermark("event_time", "3 minutes")
.groupBy(window("event_time", "1 minute"), "region")
.count()
)
Pour garantir que les agrégations sont traitées de manière incrémentielle plutôt que recalculées entièrement à chaque mise à jour, vous devez définir un watermark.
Comprendre l'état du streaming et le full refresh
L’état du streaming est incrémentiel : le pipeline construit et maintient l’état au fil des mises à jour plutôt que de le recalculer à partir de zéro à chaque fois. C’est ce qui rend le streaming avec état efficace, mais cela signifie également que si vous modifiez la logique d’une requête avec état (par exemple, en modifiant un threshold de filigrane ou en changeant les colonnes d’agrégation), l’état existant n’est plus compatible avec la nouvelle logique. Dans ce cas, vous devez effectuer un full refresh pour retraiter toutes les données historiques avec la nouvelle logique et reconstruire l’état à partir de zéro.
Un refresh complet peut également entraîner une perte de données si la source ne conserve pas les données historiques. Par exemple, une source Kafka avec une courte période de rétention peut ne disposer que des données des dernières minutes au moment du refresh, ce qui donne une table contenant beaucoup moins de données qu’auparavant. Planifiez soigneusement les modifications de la logique de query avec état, en particulier pour les Stream à haut volume où un refresh complet est coûteux ou lorsque la source a une rétention de données limitée. L’utilisation de l’ architecture en médaillon est utile car elle permet de créer des tables bronze avec une transformation minimale, et permet aux tables silver ou gold de recalculer à partir des tables bronze avec un historique complet.
Stream-Stream joins
Les jointures stream-stream nécessitent un watermark des deux côtés de la jointure et une condition de jointure limitée dans le temps. L'intervalle de temps dans la condition de jointure indique au moteur de streaming quand aucune autre correspondance n'est possible, lui permettant d'évincer l'état qui ne peut plus être mis en correspondance. Si vous omettez les watermarks ou la condition limitée dans le temps, l'état croît sans limite.
L'exemple suivant joint les événements d'impression publicitaire aux événements de clic, en exigeant que le clic se produise dans les trois minutes suivant l'impression :
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE impression_clicks AS
SELECT imp.ad_id, imp.impression_time, clk.click_time
FROM STREAM(ad_impressions)
WATERMARK impression_time DELAY OF INTERVAL 3 MINUTES AS imp
JOIN STREAM(user_clicks)
WATERMARK click_time DELAY OF INTERVAL 3 MINUTES AS clk
ON imp.ad_id = clk.ad_id
AND clk.click_time BETWEEN imp.impression_time
AND imp.impression_time + INTERVAL 3 MINUTES;
from pyspark import pipelines as dp
from pyspark.sql.functions import expr
dp.create_streaming_table("impression_clicks")
@dp.append_flow(target="impression_clicks")
def join_impressions_and_clicks():
impressions = spark.readStream.table("ad_impressions") \
.withWatermark("impression_time", "3 minutes")
clicks = spark.readStream.table("user_clicks") \
.withWatermark("click_time", "3 minutes")
return impressions.alias("imp").join(
clicks.alias("clk"),
expr("""
imp.ad_id = clk.ad_id AND
clk.click_time BETWEEN imp.impression_time AND imp.impression_time + INTERVAL 3 MINUTES
"""),
"leftOuter"
)
Lorsque vous joignez un Stream à une table statique (jointure d'instantané), l'instantané de la table statique est actualisé au start de chaque micro-lot. Cela signifie que les enregistrements de dimension arrivés tardivement ne sont pas appliqués rétroactivement aux faits qui ont déjà été traités. Si une application rétroactive est requise, utilisez une vue matérialisée ou restructurez le pipeline.
Optimiser la performance du pipeline
Appliquez ces techniques pour réduire les coûts de compute et accélérer les mises à jour du pipeline.
Pour plus d’informations, voir Vues matérialisées et Optimiser le traitement étatif avec filigranes.
Évitez les petits dossiers
Déclencher un pipeline trop fréquemment sur une source à faible volume écrit un grand nombre de petits fichiers dans le stockage cloud. Les petits fichiers dégradent la performance de lecture car chaque fichier nécessite une recherche de métadonnées et un aller-retour séparés, et les APIs de stockage cloud limitent les opérations de listage à grande échelle. Pour éviter cela, choisissez un intervalle Trigger correspondant à votre volume de données : exécutez des pipelines Trigger selon un calendrier qui permet de s’accumuler une quantité significative de données entre les mises à jour, plutôt que de manière continue.
Gérer le biais des données
Le déséquilibre des données (data skew) se produit lorsque les valeurs d'une clé de jointure ou de groupBy sont réparties de manière inégale entre les partitions, ce qui amène un petit nombre de tâches à traiter la majorité des données. Cela crée des points chauds qui augmentent le temps de mise à jour de bout en bout. Utilisez le clustering fluide pour remédier au déséquilibre dans les tables stockées. Pour l'asymétrie (skew) qui se produit pendant le calcul en cours, ajoutez un suffixe de compartiment aléatoire aux clés fortement asymétriques avant de procéder au regroupement et à l'agrégation en deux étapes.
Pour plus d'informations, consultez Utiliser le clustering fluide pour la Layout des données.
Utilisez le refresh incrémentiel pour les vues matérialisées
Lorsque vous utilisez une vue matérialisée pour une agrégation importante, le pipeline tente de l’actualiser de manière incrémentielle, en traitant uniquement les modifications en amont depuis la dernière mise à jour plutôt que de recalculer l’ensemble complet des résultats. L’actualisation incrémentielle est nettement moins coûteuse que l’exécution de la requête à partir de zéro à chaque Trigger du pipeline. Pour maximiser les chances qu’une vue matérialisée puisse être actualisée de manière incrémentielle, rédigez des requêtes d’agrégation simples et déterministes, et évitez les constructions qui empêchent le traitement incrémentiel, telles que les fonctions non déterministes.
Voir refresh incrémentielle pour les vues matérialisées.
Optimiser les jointures
Pour les jointures où l’un des côtés est une petite table de dimensions, ajoutez un indicateur de broadcast pour demander à Spark de diffuser la plus petite table vers tous les exécuteurs au lieu d’effectuer une jointure par shuffle :
- SQL
- Python
CREATE OR REFRESH MATERIALIZED VIEW enriched_orders AS
SELECT o.*, /*+ BROADCAST(p) */ p.product_name, p.category
FROM orders o
JOIN products p ON o.product_id = p.product_id;
from pyspark import pipelines as dp
from pyspark.sql.functions import broadcast
@dp.materialized_view
def enriched_orders():
orders = spark.read.table("orders")
products = spark.read.table("products")
return orders.join(broadcast(products), "product_id")
Pour les jointures de proximité de séries temporelles (par exemple, pour trouver l'événement le plus proche dans une plage de temps), utilisez une condition de jointure de plage et assurez-vous que les deux côtés ont un watermark si vous joignez des Stream, ou envisagez de pré-regrouper les événements dans des compartiments temporels avant la jointure.
Surveillez vos pipelines
Le journal des événements du pipeline est la primitive d'observabilité principale dans les pipelines. Chaque exécution de pipeline écrit des enregistrements structurés dans le journal des événements couvrant la progression de l'exécution, les résultats des attentes en matière de qualité des données, la data lineage et les détails des erreurs. Le journal des événements est une table Delta que vous pouvez interroger directement.
Pour interroger le log d’événements sans connaître le chemin de stockage sous-jacent, utilisez la fonction à valeurs de table event_log() sur un cluster partagé ou un SQL Warehouse :
SELECT * FROM event_log('<pipeline-id>')
WHERE event_type = 'flow_progress'
ORDER BY timestamp DESC
LIMIT 100;
Créez des tableaux de bord de qualité des données en interrogeant le log des événements pour obtenir les métriques d'attente. La colonne details contient une structure JSON imbriquée avec les nombres de réussites/échecs pour chaque contrainte, que vous pouvez utiliser pour suivre les tendances de qualité au fil du temps et alerter en cas de régression.
Pour les alertes basées sur les événements, utilisez des hooks d’événement pour Trigger des webhooks personnalisés ou des services de notification (tels que Slack ou PagerDuty) lorsqu’un pipeline échoue ou lorsqu’un threshold de qualité des données est dépassé. Les hooks d’événement sont des fonctions Python qui s’exécutent en réponse aux événements du pipeline.
Pour plus d’informations, consultez Superviser les pipelines, Logs des événements du pipeline et Définir une supervision personnalisée des pipelines avec des hooks d’événement.
Utiliser le calcul serverless
Databricks recommande le compute serverless pour les nouveaux pipelines. Avec le mode serverless, il n’y a aucune configuration manuelle de cluster ; Databricks gère l’infrastructure automatiquement. Les pipelines serverless utilisent une mise à l’échelle automatique améliorée qui peut monter en charge à la fois horizontalement (plus d’exécuteurs) et verticalement (taille d’exécuteur plus grande) en réponse aux exigences de la charge de travail. Les pipelines serverless utilisent toujours Unity Catalog, de sorte que la gouvernance et le suivi de la traçabilité sont intégrés par défaut.
Serverless est également requis pour l’ incremental refresh des vues matérialisées. Sur le compute classique, les vues matérialisées sont toujours entièrement recalculées, ce qui augmente les coûts de refresh. Si l’incremental refresh est importante pour votre charge de travail, utilisez le compute Serverless.
Pour une comparaison entre le Serverless et le compute classique, voir Serverless vs. compute classique pour les pipelines. Pour plus d’informations sur le serverless, voir Configurer un pipeline serverless.
Organiser les pipelines avec l'architecture en médaillon
L'architecture en médaillon organise les données en trois couches logiques (bronze, silver et gold), chacune ayant un objectif distinct. Le mappage des types de datasets de pipeline vers la bonne couche permet de clarifier les responsabilités de chaque couche et facilite la maintenance des pipelines.
- Bronze : utilisez des tables de streaming pour ingérer des données brutes à partir du stockage cloud, de bus de messages ou de sources CDC. Les tables bronze préservent les données sources brutes avec une transformation minimale, ce qui permet aux couches silver ou gold de retraiter à partir de la source dans la couche bronze si les exigences changent.
- Silver : utilisez des tables de streaming pour les transformations incrémentielles au niveau des lignes (filtrage, nettoyage et analyse). Utilisez des vues matérialisées lorsque la logique de la couche Silver implique des jointures d’enrichissement avec des tables de dimension ou des agrégations complexes qui bénéficient d’une refresh incrémentielle.
- Gold : utilisez des vues matérialisées pour pré-calculer des agrégations, des indicateurs et des résumés destinés aux tableaux de bord, aux outils de reporting et aux consommateurs en aval.
Séparez l'ingestion (bronze) et la transformation (silver et gold) dans des pipelines distincts dans la mesure du possible. Le découplage des couches vous permet de planifier, de surveiller et de dépanner chaque couche indépendamment, et une défaillance dans un pipeline de transformation n'empêche pas les nouvelles données d'arriver dans la couche bronze.
Pour plus d'informations, consultez Tables de streaming et Vues matérialisées.