Aller au contenu principal

Traitement par batch avec des flux REPLACE WHERE

Les flux REPLACE WHERE dans LakeFlow Pipelines recalculent et écrasent un sous-ensemble ciblé d'une table sans retraiter l'intégralité de 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

Définissez des flux REPLACE WHERE en SQL ou en Python.

Utilisez la clause FLOW REPLACE WHERE en ligne avec CREATE STREAMING TABLE:

SQL
CREATE STREAMING TABLE orders_enriched
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;

Alternativement, utilisez la syntaxe CREATE FLOW longue :

SQL
CREATE STREAMING TABLE orders_enriched;

CREATE FLOW orders_enriched AS
INSERT INTO orders_enriched BY NAME
REPLACE WHERE date >= date_add(current_date(), -7)
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;

Dans ces exemples, toutes les lignes des 7 derniers jours sont supprimées de orders_enriched et recalculées à l’aide de la query source. Vous n'avez pas besoin d'ajouter le prédicat à la query source. Le moteur de pipeline l’applique automatiquement lors de la lecture de la source.

remarque

BY NAME est requis en SQL. Il fait correspondre les colonnes 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 :

  • **Remplacements de prédicat** : relancez la query source du flux pour une plage de prédicats unique. 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

Outrepassez le prédicat REPLACE WHERE pour une seule mise à jour de pipeline sans modifier la définition du pipeline. Les remplacements de prédicats sont uniques, s'appliquent uniquement à la mise à jour actuelle et n'affectent pas les exécutions futures.

Exemple : Chargement historique initial

Pour effectuer un remplissage historique unique des données historiques lors de la première configuration d'un pipeline :

Python
pipeline_id = "<pipeline-id>"
overrides = [
{
"flow_name": "orders_enriched",
"predicate_override": "date BETWEEN '2020-01-01' AND '2024-12-31'",
}
]

resp = start_update_with_replace_where(
pipeline_id=pipeline_id,
replace_where_overrides=overrides,
)
print(resp)

Exemple : Corriger une colonne pour une période spécifique.

Après avoir mis à jour une définition de colonne, remplissez la modification pour une plage historique ciblée :

Python
pipeline_id = "<pipeline-id>"
overrides = [
{
"flow_name": "orders_enriched",
"predicate_override": "date >= date_add(current_date(), -30)",
}
]

resp = start_update_with_replace_where(
pipeline_id=pipeline_id,
replace_where_overrides=overrides,
refresh_selection=["orders_enriched"],
)
print(resp)

Combinez plusieurs dimensions dans un seul remplacement de prédicat :

Python
overrides = [
{
"flow_name": "orders_enriched",
"predicate_override": "date >= date_add(current_date(), -30) AND region = 'asia'",
}
]

Fonction d'aide : start_update_with_replace_where

Utilisez l'API de mise à jour de pipeline à partir d'un notebook pour soumettre des remplacements de prédicat :

Python
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import StartUpdateResponse


def start_update_with_replace_where(
pipeline_id: str,
replace_where_overrides: list[dict],
refresh_selection: list[str] = None,
) -> StartUpdateResponse:
"""Start a pipeline update with REPLACE WHERE predicate overrides."""
client = WorkspaceClient()

body = {
"pipeline_id": pipeline_id,
"cause": "JOB_TASK",
"update_cause_details": {
"job_details": {"performance_target": "PERFORMANCE"}
},
"replace_where_overrides": replace_where_overrides,
}

if refresh_selection:
body["refresh_selection"] = refresh_selection

res = client.api_client.do(
"POST",
f"/api/2.0/pipelines/{pipeline_id}/updates",
body=body,
headers={&quot;Accept&quot;: &quot;application/json&quot;, &quot;Content-Type&quot;: &quot;application/json&quot;},
)

return StartUpdateResponse.from_dict(res)

Instructions DML

Exécutez des instructions DML directement sur la table cible en dehors du pipeline pour effectuer des chargements initiaux ou des corrections, tels que le chargement à partir d'une table existante :

SQL
INSERT INTO orders_enriched
SELECT *
FROM orders_enriched_legacy
WHERE date < '2025-01-01';

Les lignes insérées via DML ne sont pas soumises au prédicat REPLACE WHERE et persistent lors des refresh planifiées, sauf si elles se situent dans la plage de prédicats d'une exécution future.

Comportement de refresh complet

Une full refresh d'un flux REPLACE WHERE réexécute la query source en utilisant uniquement le prédicat actuel. Les lignes qui ont été insérées par des remplacements de prédicats ou des instructions DML en dehors de la plage de prédicats actuelle sont supprimées définitivement.

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.

Pour éviter les actualisations complètes sur une table, définissez la propriété de table pipelines.reset.allowed sur false. Voir Référence des propriétés de pipeline.

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.

Limitations

Les flux REPLACE WHERE ont les limitations suivantes :

  • La table cible doit être créée dans le pipeline.
  • Un seul flux REPLACE WHERE est autorisé par table cible.
  • Une table ciblée par un flux REPLACE WHERE ne peut pas non plus être ciblée par un autre type de flux, tel qu'un flux CDC automatique ou un flux d'ajout.
  • Les attentes ne sont pas prises en charge sur les tables ciblées par les flux REPLACE WHERE.
  • Pour les tables de streaming autonomes, consultez les flux REPLACE WHERE pour les tables de streaming autonomes pour les différences de syntaxe et de rétroremplissage.

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 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 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. Pour corriger les lignes historiques, exécutez un remplissage ciblé en utilisant des substitutions de prédicat.

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 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;
  1. Mettez à jour la query pour ajouter uniq_users:
SQL
CREATE 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;
  1. Remplissez rétrospectivement la nouvelle métrique pour les 30 derniers jours :

    Python
    overrides = [
    {
    "flow_name": "clickstream_daily",
    "predicate_override": "event_date BETWEEN '2026-01-01' AND '2026-01-30'",
    }
    ]

    resp = start_update_with_replace_where(
    pipeline_id="<pipeline-id>",
    replace_where_overrides=overrides,
    refresh_selection=["clickstream_daily"],
    )

    Les lignes plus anciennes que la plage de remplissage 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 fenêtre courte afin que chaque refresh recalcule uniquement les 7 derniers jours pendant que vous révisez la query :

SQL
CREATE 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 fois la query finalisée, utilisez un remplacement de prédicat pour effectuer un remplissage historique unique :

Python
overrides = [
{
"flow_name": "revenue_attribution",
"predicate_override": "event_date >= date_add(current_date(), -365)",
}
]

resp = start_update_with_replace_where(
pipeline_id="<pipeline-id>",
replace_where_overrides=overrides,
refresh_selection=["revenue_attribution"],
)