Flux REPLACE WHERE pour les tables de streaming autonomes
Les flux REPLACE WHERE recalculent et écrasent un sous-ensemble ciblé d'une table de streaming autonome sans retraiter tout l'historique de votre table. Ils gèrent les données arrivées en retard, le retraitement en amont, l’évolution des schémas et les rétroremplissages.
Avec un flux REPLACE WHERE, vous définissez un prédicat sur la table cible. Toutes les lignes correspondant au prédicat sont supprimées et remplacées en réévaluant la source query pour cette même plage de prédicats. Les lignes qui ne correspondent pas au prédicat restent inchangées.
Exigences
Les flux REPLACE WHERE ont les exigences suivantes :
- Databricks recommande Unity Catalog et le compute serverless. Incremental refresh est uniquement pris en charge sur le compute Serverless.
Quand utiliser les flux REPLACE WHERE
Utilisez les flux REPLACE WHERE pour les scénarios suivants :
- **Traitement par batch incrémentiel sans sémantique de streaming** : Traitez de nouvelles lignes par batchs sans gérer les concepts de streaming tels que les filigranes.
- **Retraitement sélectif** : Recalculez uniquement les lignes qui correspondent à un prédicat tout en laissant toutes les autres lignes intactes.
- Scénarios au-delà des capacités des vues matérialisées standard :
- Tables cibles avec une conservation plus longue que la source
- Prévention du recalcul lors de la modification d'une table de dimension
- Évolution des schémas sans recalculer l'historique complet
Créer un flux REPLACE WHERE
Utilisez la clause FLOW REPLACE WHERE en ligne avec CREATE OR REFRESH STREAMING TABLE:
CREATE OR REFRESH STREAMING TABLE orders_enriched
SCHEDULE EVERY 1 DAY
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
SELECT
o.order_id,
o.date,
o.region,
p.product_name,
o.qty,
o.price
FROM orders_fct o
JOIN product_dim p
ON o.product_id = p.product_id;
Lors du refresh, toutes les lignes de la table cible qui correspondent au prédicat sont supprimées, la query source est recalculée pour la même plage de prédicats et les nouveaux résultats sont insérés. Dans cet exemple, toutes les lignes des 7 derniers jours sont supprimées de orders_enriched et recalculées à l'aide de la source query.
Vous n'avez pas besoin d'ajouter le prédicat à la query source. Le moteur de pipeline l'applique automatiquement lors de la lecture depuis la source.
BY NAME est obligatoire. Cela garantit que les colonnes sont mises en correspondance par nom plutôt que par position.
Remplir les données historiques
Pour écrire des lignes historiques ou corrigées dans la table cible en dehors des refresh planifiés, choisissez entre deux mécanismes en fonction de l'emplacement des données historiques :
- Écrasements de prédicats : Réexécutez la query source du flux pour une plage de prédicats unique en utilisant
REFRESH STREAMING TABLE ... WHERE. Utilisez lorsque les données historiques proviennent de la même source que les données incrémentielles. - Instructions DML : Insérez directement dans la table cible, en contournant le flux. Utilisez lorsque les données historiques se trouvent dans une source différente de celle des données incrémentielles.
Remplacements de prédicat
Dérogiez le prédicat REPLACE WHERE du flux pour un seul refresh sans modifier la définition de la table. Les dérogations de prédicat sont uniques, s'appliquent uniquement à ce refresh et n'affectent pas les refresh planifiés ultérieurs.
Utilisez la clause WHERE avec REFRESH STREAMING TABLE:
REFRESH STREAMING TABLE orders_enriched WHERE date BETWEEN '2020-01-01' AND '2024-12-31';
Le flux supprime uniquement les lignes correspondant au prédicat d'écrasement et les recalcule à partir de la source. Le prédicat statique du flux est inchangé, de sorte que la prochaine scheduled refresh utilise le prédicat d'origine.
Vous pouvez combiner le remplacement avec SYNC ou ASYNC. Par exemple, pour start le refresh en arrière-plan et revenir immédiatement avec un Link vers la mise à jour du pipeline :
REFRESH STREAMING TABLE orders_enriched WHERE date >= date_add(current_date(), -30) ASYNC;
La clause WHERE n'est prise en charge que sur une table de streaming créée avec une clause FLOW REPLACE WHERE. Le spécifier sur toute autre table streaming renvoie une erreur. Voir REFRESH (MATERIALIZED VIEW ou STREAMING TABLE).
Instructions DML
Exécutez les instructions DML directement sur la table cible pour charger les lignes provenant d'une source différente de celle du flux :
INSERT INTO orders_enriched
SELECT *
FROM orders_enriched_legacy
WHERE date < '2025-01-01';
Comportement de refresh complet
Un refresh complet d'un flux REPLACE WHERE réexécute la query source en utilisant uniquement le prédicat défini du flux. Il supprime définitivement chaque ligne de la cible qui ne correspond pas à ce prédicat, y compris les lignes que les écrasements de prédicats ou les instructions DML ont précédemment insérées en dehors de la plage de prédicats définie.
Un refresh complet efface toutes les données existantes et réexécute le flux en utilisant uniquement son prédicat défini. Si un pipeline a été exécuté pendant un an avec un prédicat de 7 jours, une refresh complète fait en sorte que la table ne contient que les données des 7 derniers jours. Toutes les lignes plus anciennes sont supprimées définitivement.
REFRESH STREAMING TABLE orders_enriched FULL;
Pour éviter les actualisations complètes d'une table, définissez la propriété de la table pipelines.reset.allowed sur false:
CREATE OR REFRESH STREAMING TABLE orders_enriched
TBLPROPERTIES (pipelines.reset.allowed = 'false')
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
...
refresh incrémentielle
Les flux REPLACE WHERE utilisent une refresh incrémentielle lorsque cela est possible, retraitant uniquement les données sources qui ont changé depuis la dernière refresh plutôt que de recalculer toute la fenêtre de remplacement. Le refresh incrémentiel nécessite un compute Serverless.
Quand l'incremental refresh s'applique
Tous les éléments suivants doivent être vrais :
- Le pipeline s'exécute sur un compute serverless.
- La forme de la query est prise en charge. Voir refresh incrémentiel pour l’ensemble d’opérateurs pris en charge.
- Le prédicat fait référence aux colonnes de base d'une table source. Les prédicats sur les valeurs dérivées, telles que les sorties de fonctions d'agrégation ou de fenêtre, ne peuvent pas être poussés vers une source, ce qui désactive le refresh incrémentiel.
- Aucun DML externe n'a modifié de lignes dans la fenêtre de remplacement actuelle. Les DML qui modifient des lignes en dehors de la fenêtre actuelle ne sont pas affectés.
- **La fenêtre de remplacement actuelle n’inclut pas les lignes que le prédicat précédent a exclues.** Si vous élargissez le prédicat pour couvrir une plage non traitée précédemment, ce refresh unique revient à une recomputation complète. Les prochains refresh sont de nouveau éligibles au refresh incrémentiel.
- Le prédicat est déterministe. Les prédicats utilisant des fonctions non déterministes telles que
rand()désactivent le refresh incrémentiel. Les fonctions temporelles telles quecurrent_date()sont autorisées.
La première refresh de tout flux est toujours un calcul complet. Si une condition n'est pas remplie, ce refresh revient à un recalcul complet de la fenêtre de remplacement actuelle.
Meilleures pratiques pour le refresh incrémentiel
Suivez ces directives afin que les flux REPLACE WHERE restent éligibles au refresh incrémentiel.
Utilisez une borne inférieure mobile
Les prédicats avec une borne inférieure mobile restent éligibles pour un refresh incrémentiel indéfiniment.
FLOW REPLACE WHERE date >= date_add(current_date(), -7)
Une limite supérieure mobile, telle que date BETWEEN date_add(current_date(), -7) AND current_date(), peut déplacer la fenêtre pour inclure des lignes précédemment exclues, ce qui déclenche un fallback unique vers un recalcul complet.
Inclure la colonne de prédicat dans GROUP BY
Lors de l'agrégation, incluez la colonne de prédicat dans GROUP BY afin que le moteur puisse pousser le prédicat sous l'agrégation.
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
SELECT date, region, SUM(amount) AS total
FROM sales
GROUP BY date, region;
Si la colonne de prédicat manque à GROUP BY, le prédicat ne peut pas être poussé sous l'agrégation et la source est analysée en intégralité.
Inclure la colonne de prédicat dans les clés de jointure
Incluez la colonne de prédicat dans la condition de jointure afin que le moteur puisse élaguer toutes les sources jointes.
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
SELECT f.date, f.user_id, d.region, f.revenue
FROM fact f
JOIN dim d ON f.date = d.date AND f.user_id = d.user_id;
Si une table jointe n'expose pas la colonne de prédicat, cette table est entièrement analysée à chaque refresh.
Diagnostiquer le fallback vers un recalcul complet
Lorsqu'un refresh rebascule vers un recalcul complet, la raison est signalée dans l'événement planning_information pour le flux. Consultez Surveiller les Logs d'événements de pipeline. Le tableau suivant répertorie les raisons signalées dans l'événement :
Raison | Signification |
|---|---|
| Un DML externe a modifié des lignes dans la fenêtre de remplacement actuelle. |
| Le prédicat utilise des expressions non déterministes. |
| Le refresh précédent a utilisé un prédicat non déterministe. |
| Le prédicat ne peut pas être transféré vers une source, la fenêtre actuelle inclut des lignes non traitées par le prédicat précédent, ou l'exécution utilise une substitution de prédicat. |
Exemples
Les exemples suivants illustrent les modèles de flux REPLACE WHERE courants.
Exemple 1 : Conserver les agrégats historiques d'une source à rétention limitée.
Cet exemple conserve les agrégats quotidiens indéfiniment, même après que les données brutes ont expiré de la table source (conservation de 3 jours) :
CREATE OR REFRESH STREAMING TABLE events_agg
FLOW REPLACE WHERE date >= date_add(current_date(), -3) BY NAME
SELECT
date,
key,
SUM(val) AS agg
FROM events_raw
GROUP BY ALL;
Exemple 2 : Empêcher le recalcul en cas de modification d'une table de dimension
Cet exemple conserve les lignes de faits historiques inchangées lorsque les attributs de dimension changent :
CREATE OR REFRESH STREAMING TABLE fact_dim_join
FLOW REPLACE WHERE f.date >= date_add(current_date(), -1) BY NAME
SELECT
f.date,
f.user_id,
d.region,
f.revenue
FROM fact_table f
JOIN dim_users d
ON f.user_id = d.user_id;
Si la région d'un utilisateur change, seules les lignes récentes sont recalculées. Les lignes historiques conservent la valeur de la région au moment où elles ont été écrites.
Exemple 3 : Ajouter une nouvelle métrique sans recalculer l'historique complet
Cet exemple montre comment faire évoluer une définition de table et ne remplir qu'une plage ciblée :
-
Définissez la table initiale :
SQLCREATE OR REFRESH STREAMING TABLE clickstream_daily
FLOW REPLACE WHERE event_date >= date_add(current_date(), -7) BY NAME
SELECT
event_date,
page_id,
COUNT(*) AS clicks
FROM clickstream_raw
GROUP BY ALL; -
Mettez à jour la query pour ajouter
uniq_users:SQLCREATE OR REFRESH STREAMING TABLE clickstream_daily
FLOW REPLACE WHERE event_date >= date_add(current_date(), -7) BY NAME
SELECT
event_date,
page_id,
COUNT(*) AS clicks,
COUNT(DISTINCT user_id) AS uniq_users
FROM clickstream_raw
GROUP BY ALL;Les lignes plus anciennes que la fenêtre de 7 jours contiennent
NULLpouruniq_users.
Exemple 4 : Itérer sur une petite fenêtre avant de reconstituer l'historique complet
Cet exemple montre comment valider la logique de query sur une petite fenêtre de données avant de traiter la plage historique complète.
Start avec une courte période pour valider les métriques et itérer sur la logique métier avec des coûts de compute réduits :
CREATE OR REFRESH STREAMING TABLE revenue_attribution
FLOW REPLACE WHERE event_date >= date_add(current_date(), -7) BY NAME
SELECT
event_date,
campaign_id,
SUM(revenue) AS total_revenue
FROM marketing_events
GROUP BY ALL;
Une fenêtre courte ne recalcule que les 7 derniers jours à chaque refresh, veuillez donc réviser la requête autant de fois que nécessaire avant de la commit pour une exécution historique complète.
Une fois la query finalisée, utilisez DML pour remplir la plage historique complète :
INSERT INTO revenue_attribution
SELECT
event_date,
campaign_id,
SUM(revenue) AS total_revenue
FROM marketing_events
WHERE event_date < date_add(current_date(), -7)
GROUP BY ALL;