Bonnes pratiques pour les LakeFlow Pipelines
Appliquez ces modèles recommandés lors de la conception, de la construction et de l'exploitation des pipelines, que vous démarriez un nouveau pipeline ou que vous en amélioriez un existant.
Choisissez le bon type de dataset
Les pipelines proposent trois types de dataset : les tables de streaming, les vues matérialisées et les vues temporaires. Choisir le bon type pour chaque couche de votre pipeline permet d’éviter des coûts de compute inutiles et rend votre code plus facile à comprendre.
Les tables de streaming sont le bon choix 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 workloads en mode ajout seul, les données à volume élevé et le traitement événementiel à partir du stockage cloud ou des bus de messages.
Les vues matérialisées sont le bon choix pour les transformations complexes et les queries analytiques. Leurs résultats sont précalculés et mis à jour à l'aide d'une incremental refresh, de sorte que les requêtes les concernant sont rapides. Vous ne pouvez pas modifier directement les données dans une vue matérialisée ; la définition de la query contrôle la sortie.
Les vues temporaires sont des vues à l'échelle du pipeline qui organisent votre logique de transformation 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écapitule quand utiliser chaque type :
Cas d'usage | Type recommandé | Raison |
|---|---|---|
Ingestion depuis un stockage cloud ou un bus de messages | Table de streaming | Traite chaque enregistrement une seule fois ; gère les charges de travail à volume élevé et en ajout uniquement. |
CDC Stream (insertions, mises à jour, suppressions) | Table de streaming | Utilisé comme cible de |
Agrégations et jointures complexes | Vue matérialisée | Actualisé de manière incrémentielle ; évite un nouveau calcul complet à chaque mise à jour. |
Accélération de la query du tableau de bord | Vue matérialisée | Les résultats précalculés rendent les queries plus rapides que les tables brutes. |
Transformations intermédiaires (pas de lecteurs en aval) | Vue temporaire | Organise la logique du pipeline sans occasionner de coûts de stockage. |
Pour plus d'informations, consultez les tables en streaming, les vues matérialisées et Qu'est-ce que Lakeflow Pipelines ?
Utilisez la CDC déclarative au lieu de MERGE impérative
La mise en œuvre de la capture de données modifiées (CDC) avec des instructions SQL MERGE impératives nécessite un code personnalisé significatif pour gérer correctement l'ordonnancement 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 résultant 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 le classement, la déduplication, les événements désordonnés et l'évolution des schémas de manière déclarative. Vous décrivez la forme du flux de modifications et la table cible, et le pipeline gère le reste. AUTO CDC prend en charge les SCD de type 1 (écrasement) et les SCD de type 2 (conservation de l'historique).
Pour plus d'informations, consultez Capture de données modifiées et instantanés et Les APIs AUTO CDC : Simplifier la capture de données modifiées avec des pipelines.
Assurer la qualité des données avec des attentes
Les attentes sont des expressions SQL vrai/faux appliquées à chaque ligne passant par un dataset. Lorsqu'une ligne ne remplit pas la condition, le pipeline réagit conformément à la politique de violation que vous avez configurée. Les attentes émettent des métriques vers le journal des événements de pipeline quelle que soit la politique, afin que vous puissiez suivre les tendances de qualité des données au fil du temps.
Choisissez une politique de violation
Trois stratégies de violation sont disponibles. Choisissez celui qui correspond à votre tolérance aux mauvaises données :
- warn (default) : les enregistrements non valides sont écrits dans la table cible et signalés dans les métriques. Utilisez cette politique lorsque vous avez besoin de capturer toutes les données mais souhaitez avoir une visibilité sur les problèmes de qualité.
- drop : Les enregistrements qui ne sont pas valides sont ignorés avant l'écriture. Utilisez ceci lorsque des lignes incorrectes sont attendues et ne doivent pas se propager en aval.
- échec : La mise à jour du pipeline s'arrête au premier enregistrement non valide. Utilisez ceci pour les données critiques où tout enregistrement incorrect indique un grave problème 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 abandonnés à des fins d'investigation plutôt que de les abandonner 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 défectueuses 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étrez vos pipelines
Les pipelines disposent de paramètres de catalogue et de schéma default, de sorte que le code qui lit et écrit au sein du même catalogue et schéma fonctionne dans différents environnements sans aucun paramètre. Toutefois, si votre pipeline doit référencer un second catalogue ou schéma (par exemple, lire à partir d'un catalogue source partagé qui diffère entre le développement et la production), évitez de coder en dur ces noms directement dans votre code source. Au lieu de cela, définissez-les comme des paramètres de configuration de pipeline (paires clé-valeur définies dans les paramètres du pipeline) et référencez-les dans votre code. Cela permet à une base de code unique de s'exécuter correctement dans différents environnements en échangeant 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 les paramètres avec les pipelines.
Choisissez entre le mode de pipeline déclenché et 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 un planning (horaire, quotidien ou à la demande) et ne nécessitent pas une actualisation des données en moins d'une minute.
Le mode continu maintient le cluster en cours d’exécution et traite les nouvelles données au fur et à mesure qu’elles arrivent. Cela convient uniquement lorsque votre cas d’utilisation nécessite une latence de l’ordre de quelques secondes à quelques minutes. Étant donné que le mode continu nécessite un cluster toujours actif, il est beaucoup plus coûteux que le mode Trigger.
Le mode temps réel s'appuie sur le mode continu pour atteindre une latence de l'ordre de la milliseconde pour les charges de travail opérationnelles telles que la détection de fraude ou la personnalisation en temps réel. Cela nécessite une configuration et une planification compute supplémentaires. Consultez Utiliser le mode temps réel dans les LakeFlow Pipelines.
Pour plus d’informations, consultez Mode de pipeline Trigger ou continu et Configurer les pipelines.
Utiliser le clustering liquide pour le Layout des données
Le clustering liquide remplace le partitionnement statique et ZORDER pour l'optimisation du layout des données dans les tables Delta. Le partitionnement statique vous oblige à sélectionner les colonnes de partition et à réorganiser les données à l'avance, ce qui peut entraîner une asymétrie des données pour les valeurs inégalement réparties. Le clustering liquide est auto-ajustant, résistant aux asymétries et incrémentiel, ne réécrivant que les données qui nécessitent une réorganisation à chaque exécution.
Modifiez les colonnes de clustering à tout moment sans réécrire l'intégralité de la table à mesure que les schémas 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 requêtes. 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 clustering, 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, consultez Tables de streaming et Utiliser le clustering liquide pour les tables.
Gérez les pipelines avec CI/CD et les Declarative Automation Bundles
Contrôlez la version de votre code source de pipeline et utilisez les Declarative Automation Bundles pour gérer les déploiements dans tous les environnements.
Pour plus d'informations, consultez Créer un pipeline contrôlé par le code source, Convertir un pipeline en projet de bundle et Utiliser des paramètres avec des pipelines.
Stocker le code de pipeline dans un système de 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 du projet complet vous donne un historique complet des modifications, facilite la collaboration et vous permet de valider les modifications dans un environnement de développement avant de les promouvoir en production.
Databricks recommande les Declarative Automation Bundles pour la gestion de ce workflow. Un bundle définit votre configuration de pipeline en YAML aux côtés de votre code source, et le databricks bundle CLI vous permet de valider, déployer et exécuter des pipelines depuis votre terminal ou un système CI/CD.
Utiliser les cibles de bundle pour l’isolation de l’environnement
Les bundles permettent de multiples *cibles* (par exemple,,,), dev staging``prodchacune avec son propre ensemble de remplacements pour les noms de catalogue, les politiques de clusters, les adresses de notification et d'autres paramètres. Combinez les cibles de bundle avec les parameters de pipeline pour injecter les valeurs spécifiques à l'environnement correctes au moment du déploiement, en gardant votre code source exempt de constantes d'environnement.
Un workflow typique se présente comme suit :
- Un développeur travaille sur une feature Branch, en déployant vers 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 de CI déploie en production avec
databricks bundle deploy --target prod.
Bonnes pratiques de streaming
Utilisez ces modèles pour gérer l'état, contrôler les données tardives et maintenir la fiabilité des pipelines de streaming.
Pour plus d'information, consultez Optimiser le traitement avec état avec les filigranes, 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.
Utilisez des filigranes pour les opérations avec état
Les filigranes lient l’état que le pipeline conserve en mémoire pendant les opérations de streaming avec état telles que les agrégations fenêtrées et la déduplication. Sans un filigrane, l'état croît de manière illimitée à mesure que le pipeline accumule des données pour chaque clé possible, ce qui finit par provoquer des erreurs de mémoire insuffisante sur les pipelines de longue durée.
Un filigrane spécifie une colonne de timestamp et un threshold de tolérance pour les données tardives. Les enregistrements qui arrivent après le dépassement du threshold sont supprimés. Choisissez un threshold qui équilibre votre tolérance aux données tardives et le coût mémoire du maintien de cet état ouvert.
L'exemple suivant calcule une agrégation de fenêtre glissante 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 d'être entièrement recalculées à chaque mise à jour, vous devez définir un watermark.
Comprendre l'état du streaming et le full refresh
L'état de streaming est incrémentiel : le pipeline construit et maintient l'état à travers les 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 query avec état (par exemple, en changeant un threshold de watermark ou en modifiant des 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 n'avoir que les dernières minutes de données disponibles au moment du refresh, ce qui entraîne une table qui contient beaucoup moins de données qu'auparavant. Planifiez attentivement les modifications de la logique des requêtes avec état, en particulier pour les Stream à volume élevé 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 aide en créant des tables bronze avec une transformation minimale, et permet aux tables silver ou gold de se recalculer à partir des tables bronze avec l'historique complet.
Jointures stream-stream
Les jointures Stream-Stream nécessitent un filigrane des deux côtés de la jointure et une condition de jointure bornée dans le temps. L'intervalle de temps dans la condition de jointure indique au moteur de streaming qu'aucune autre correspondance n'est possible, lui permettant ainsi de supprimer les états qui ne peuvent plus être mis en correspondance. Si vous omettez soit les filigranes, soit la condition bornée dans le temps, l'état croît sans limite.
L'exemple suivant joint les événements d'impression d'annonce avec les événements de clic, exigeant que le clic ait lieu 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 (une 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 en retard 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 les performances du pipeline
Appliquez ces techniques pour réduire les coûts de compute et accélérer les mises à jour des pipelines.
Pour plus d'informations, consultez Vues matérialisées et Optimiser le traitement avec état à l'aide de filigranes.
Éviter les petits fichiers
Le fait de 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 les performances de lecture car chaque fichier nécessite une recherche de métadonnées et un aller-retour d'E/S distincts, et les APIs de stockage cloud limitent les opérations de listage à grande échelle. Pour éviter cela, choisissez un intervalle de trigger qui correspond à votre volume de données : exécutez les pipelines déclenchés selon un calendrier qui permet à une quantité significative de données de s'accumuler entre les mises à jour, plutôt que de manière continue.
Gestion de l'asymétrie des données
L'asymétrie des données se produit lorsque les valeurs d'une clé de jointure ou de regroupement (groupBy) sont inégalement réparties entre les partitions, ce qui entraîne qu'un petit nombre de tâches traite 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 liquide pour gérer l'asymétrie dans les tables stockées. Pour l'asymétrie qui se produit lors des calculs en cours, salez les clés fortement asymétriques en ajoutant un suffixe de compartiment aléatoire avant le regroupement et l'agrégation en deux étapes.
Pour plus d'informations, consultez Utiliser le clustering liquide pour le Layout des données.
Utilisez l'incremental refresh pour les vues matérialisées.
Lorsque vous utilisez une vue matérialisée pour une agrégation importante, le pipeline tente de la refresh de manière incrémentielle, en traitant uniquement les changements en amont depuis la dernière mise à jour plutôt que de recalculer l'ensemble du résultat. Le refresh incrémentiel est significativement moins coûteux que de réexécuter la requête à partir de zéro à chaque trigger de pipeline. Pour maximiser la chance qu'une vue matérialisée puisse être refresh de manière incrémentielle, écrivez 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 Incremental refresh pour les vues matérialisées.
Optimiser les jonctions
Pour les jointures où un côté est une petite table de dimension, ajoutez une indication de diffusion pour indiquer à Spark de diffuser la plus petite table à tous les exécuteurs au lieu d'effectuer une jointure par brassage :
- 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 chronologiques (par exemple, la recherche de 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 filigrane si vous joignez des Stream, ou envisagez de pré-compartimenter les événements en intervalles de temps avant de joindre.
Supervisez vos pipelines
Le journal des événements du pipeline est la primitive principale d’observabilité dans les pipelines. Chaque exécution du pipeline écrit des enregistrements structurés dans le journal des événements couvrant la progression de l’exécution, les résultats des attentes de qualité des données, le data lineage et les détails des erreurs. Le journal des événements est une table Delta que vous pouvez interroger directement.
Pour query le Log d’événements sans connaître le chemin de stockage sous-jacent, utilisez la fonction 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 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 sur les régressions.
Pour l’alerte pilotée par les événements, utilisez des crochets d’événement pour Trigger des webhooks personnalisés ou des services de notification (tels que Slack ou PagerDuty) lorsqu’un pipeline échoue ou qu’un threshold de qualité des données est franchi. Les crochets d’événement sont des fonctions Python qui s’exécutent en réponse aux événements de pipeline.
Pour plus d'information, consultez Monitorer les pipelines, Logs d'événements du pipeline et Définir le monitoring personnalisé des pipelines avec des hooks d'événements.
Utiliser le compute Serverless
Databricks recommande le compute Serverless pour les nouveaux pipelines. Avec Serverless, il n'y a pas de configuration manuelle de clusters ; Databricks gère automatiquement l'infrastructure. Les pipelines Serverless utilisent un autoscaling amélioré qui peut monter en charge horizontalement (plus d'executors) et verticalement (taille d'executor 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 default.
Pour plus d'informations, consultez 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 dataset de pipeline à 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 de stockage cloud, de bus de messages ou de sources CDC. Les tables Bronze préservent les données brutes sources avec une transformation minimale, ce qui permet aux couches Silver ou Gold de retraiter les données depuis 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 argent implique des jointures d’enrichissement par rapport aux tables de dimensions ou des agrégations complexes qui bénéficient d’un refresh incémentiel.
- **Gold** : Utilisez les vues matérialisées pour pré-calculer les agrégations, les métriques et les 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) en pipelines distincts chaque fois que 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 transformations n'empêche pas les nouvelles données d'atterrir dans la couche Bronze.
Pour plus d'information, consultez les tables de streaming et les vues matérialisées.