Aller au contenu principal

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:

SQL
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.

remarque

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:

SQL
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 :

SQL
REFRESH STREAMING TABLE orders_enriched WHERE date >= date_add(current_date(), -30) ASYNC;
remarque

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 :

SQL
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.

attention

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.

SQL
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:

SQL
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 que current_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.

SQL
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.

SQL
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.

SQL
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

EXTERNAL_CHANGE_IN_REPLACE_WINDOW

Un DML externe a modifié des lignes dans la fenêtre de remplacement actuelle.

REPLACE_WHERE_NOT_DETERMINISTIC

Le prédicat utilise des expressions non déterministes.

PRIOR_REPLACE_WHERE_NOT_DETERMINISTIC

Le refresh précédent a utilisé un prédicat non déterministe.

UNSUPPORTED_REPLACE_WHERE_PREDICATE

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.

Raison

Signification

EXTERNAL_CHANGE_IN_REPLACE_WINDOW

Un DML externe a modifié des lignes dans la fenêtre de remplacement actuelle.

REPLACE_WHERE_NOT_DETERMINISTIC

Le prédicat utilise des expressions non déterministes.

PRIOR_REPLACE_WHERE_NOT_DETERMINISTIC

Le refresh précédent a utilisé un prédicat non déterministe.

UNSUPPORTED_REPLACE_WHERE_PREDICATE

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) :

SQL
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 :

SQL
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 :

  1. Définissez la table initiale :

    SQL
    CREATE 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;
  2. Mettez à jour la query pour ajouter uniq_users:

    SQL
    CREATE 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 NULL pour uniq_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 :

SQL
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 :

SQL
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;