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.
- SQL
- Python
Utilisez la clause FLOW REPLACE WHERE en ligne avec CREATE STREAMING TABLE:
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 :
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;
En Python, la table et le flux sont définis dans une seule instruction. Le flux hérite du même nom que la table :
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("date") >= F.date_sub(F.current_date(), 7)
)
def orders_enriched():
orders_fct = spark.read.table("orders_fct").select("date", "order_id", "region", "qty", "price")
product_dim = spark.read.table("product_dim")
return orders_fct.join(product_dim, "product_id")
Le paramètre replace_where accepte soit une expression de colonne PySpark, soit un prédicat de chaîne.
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.
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 :
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 :
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 :
overrides = [
{
"flow_name": "orders_enriched",
"predicate_override": "date >= date_add(current_date(), -30) AND region = 'asia'",
}
]
Fonction d'aide : start_update_with_replace_where
start_update_with_replace_whereUtilisez l'API de mise à jour de pipeline à partir d'un notebook pour soumettre des remplacements de prédicat :
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={"Accept": "application/json", "Content-Type": "application/json"},
)
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 :
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.
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 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. |
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
- Python
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;
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("date") >= F.date_sub(F.current_date(), 3)
)
def events_agg():
return (
spark.read.table("events_raw")
.groupBy("date", "key")
.agg(F.sum("val").alias("agg"))
)
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
- Python
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;
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("date") >= F.date_sub(F.current_date(), 1)
)
def fact_dim_join():
fact_table = spark.read.table("fact_table").alias("f")
dim_users = spark.read.table("dim_users").alias("d")
return (
fact_table.join(dim_users, col("f.user_id") == col("d.user_id"))
.select(
col("f.date"),
col("f.user_id"),
col("d.region"),
col("f.revenue"),
)
)
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 :
- Définissez la table initiale :
- SQL
- Python
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;
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("event_date") >= F.date_sub(F.current_date(), 7)
)
def clickstream_daily():
return (
spark.read.table("clickstream_raw")
.groupBy("event_date", "page_id")
.agg(F.count("*").alias("clicks"))
)
- Mettez à jour la query pour ajouter
uniq_users:
- SQL
- Python
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;
@dp.table(
replace_where=col("event_date") >= F.date_sub(F.current_date(), 7)
)
def clickstream_daily():
return (
spark.read.table("clickstream_raw")
.groupBy("event_date", "page_id")
.agg(
F.count("*").alias("clicks"),
F.countDistinct("user_id").alias("uniq_users"),
)
)
-
Remplissez rétrospectivement la nouvelle métrique pour les 30 derniers jours :
Pythonoverrides = [
{
"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
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 fenêtre courte afin que chaque refresh recalcule uniquement les 7 derniers jours pendant que vous révisez la query :
- SQL
- Python
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;
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("event_date") >= F.date_sub(F.current_date(), 7)
)
def revenue_attribution():
return (
spark.read.table("marketing_events")
.groupBy("event_date", "campaign_id")
.agg(F.sum("revenue").alias("total_revenue"))
)
Une fois la query finalisée, utilisez un remplacement de prédicat pour effectuer un remplissage historique unique :
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"],
)