Remplacement partiel d'instantané avec les flux REPLACE USING
Bêta
Cette fonctionnalité est en version Bêta.
Un flux REPLACE USING maintient une table cible synchronisée avec une source de streaming : il remplace toutes les lignes qui correspondent aux colonnes clés spécifiées et laisse toutes les autres données inchangées.
Une colonne SEQUENCE BY ordonne les mises à jour afin que le résultat soit correct, même lorsque les mises à jour arrivent dans le désordre. Pour chaque clé, la séquence la plus élevée l'emporte, et une ligne avec une séquence inférieure ne remplace jamais une ligne avec une séquence supérieure déjà présente dans la cible. Les lignes qui partagent la même clé et la même séquence sont ajoutées plutôt que remplacées.
Comment fonctionne REPLACE USING
Considérez une table d’événements qui contient des événements de clic et de conversion pour deux régions, séquencés par seq:
region_id | type_de_périphérique | event_type | seq |
|---|---|---|---|
1 | iOS | cliquez | 1 |
1 | Android | conversion | 1 |
2 | iOS | cliquez | 1 |
2 | ordinateur de bureau | cliquez | 1 |
Un flux REPLACE USING (region_id) SEQUENCE BY seq reçoit ces mises à jour pour les régions 1 et 3. La région 2 ne reçoit aucune mise à jour :
region_id | type_de_périphérique | event_type | seq |
|---|---|---|---|
1 | iOS | cliquez | 2 |
1 | Android | conversion | 2 |
1 | ordinateur de bureau | cliquez | 2 |
3 | iOS | cliquez | 1 |
3 | ordinateur de bureau | cliquez | 2 |
La cible devient :
region_id | type_de_périphérique | event_type | seq | Résultat |
|---|---|---|---|---|
1 | iOS | cliquez | 2 | Remplacé, car la séquence 2 est supérieure à la séquence 1 |
1 | Android | conversion | 2 | Remplacé, car la séquence 2 est supérieure à la séquence 1 |
1 | ordinateur de bureau | cliquez | 2 | Remplacé, car la séquence 2 est supérieure à la séquence 1 |
2 | iOS | cliquez | 1 | Inchangé, car la clé n'est pas présente dans cette mise à jour |
2 | ordinateur de bureau | cliquez | 1 | Inchangé, car la clé n'est pas présente dans cette mise à jour |
3 | ordinateur de bureau | cliquez | 2 | Ajouté. La ligne seq 1 pour la région 3 n’est pas ajoutée, car seule la séquence la plus élevée pour une clé est appliquée. |
Exigences
Les flux REPLACE USING ont les exigences suivantes :
- REPLACE USING flows s’exécute sur Databricks Runtime 18.2 et versions supérieures, sur un compute classique ou serverless. Databricks recommande Unity Catalog.
- La source doit être une source de streaming. REPLACE USING rejette toute source qui n'est pas une source de streaming.
- Vous devez spécifier au moins une colonne de clé et exactement une colonne
SEQUENCE BY.
Quand utiliser les flux REPLACE USING
Les LakeFlow Pipelines proposent trois flux qui écrasent les lignes existantes. Faites votre choix en fonction de l'apparence de votre source et de la manière dont elle identifie les lignes à remplacer :
- Utilisez REPLACE USING lorsque votre source est une série d’instantanés partiels indexés par colonne. REPLACE USING écrase uniquement les données qui ont une correspondance dans les données entrantes, laissant toutes les autres données intactes. Elle ne nécessite pas de clé primaire.
- Utilisez AUTO CDC lorsque votre source est un flux de capture de données de changement (CDC) avec des opérations explicites d' insertion , de mise à jour et de suppression , ou si vous avez besoin d'un historique de dimension à évolution lente (SCD) de type 2 . AUTO CDC nécessite également une véritable clé primaire. Voir Les AUTO CDC APIs : simplifiez la capture des données de changement avec les pipelines.
- Utilisez REPLACE WHERE lorsque votre source est un snapshot et que vous souhaitez recalculer et remplacer une plage de la table cible sélectionnée par un prédicat, par exemple les 7 derniers jours, en tant qu'opération de batch. Il ne nécessite pas de clé primaire. Voir Traitement par batch avec les flux REPLACE WHERE.
Créer un flux REPLACE USING
Définissez les flux REPLACE USING en SQL ou en Python.
- SQL
- Python
Utilisez la clause FLOW REPLACE USING en ligne avec CREATE STREAMING TABLE:
CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);
Sinon, utilisez la syntaxe longue CREATE FLOW :
CREATE STREAMING TABLE payments_current;
CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);
BY NAME est requis en SQL. Il fait correspondre les colonnes par nom plutôt que par position.
Déclarez la table et le flux ensemble avec @dp.table:
from pyspark import pipelines as dp
@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
return spark.readStream.table("samples.wanderbricks.payments")
Alternativement, ciblez une table de streaming existante avec @dp.replace_flow:
from pyspark import pipelines as dp
dp.create_streaming_table("payments_current")
@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
return spark.readStream.table("samples.wanderbricks.payments")
replace_using est une liste de colonnes clés. sequence_by est un nom de colonne ou une expression Column, et est requis chaque fois que replace_using est défini.
Séquençage et données désordonnées
La colonne SEQUENCE BY rend le résultat indépendant de l'ordre dans lequel les mises à jour arrivent. Une ligne est appliquée à une clé uniquement si sa séquence est supérieure à la séquence déjà stockée pour cette clé ; ainsi, une ligne tardive ou rejouée qui est plus ancienne que la valeur actuelle est ignorée. Les clés absentes d'une mise à jour ne sont pas modifiées.
Suivez ces pratiques pour que le remplacement se comporte de manière prévisible :
Pratique | Motif |
|---|---|
Utilisez une séquence qui augmente strictement par version de clé, telle qu'un Timestamp, un numéro de version ou un décalage de log. | Deux lignes ayant la même clé et la même séquence sont toutes deux conservées, ce qui entraîne des lignes dupliquées pour cette clé. |
Utilisez une séquence non nulle. | Une séquence nulle peut entraîner un comportement non défini. |
Attentes
Les flux REPLACE USING prennent en charge les attentes. warn et fail se comportent comme sur les autres flux : warn conserve les lignes en violation et enregistre la violation, et fail arrête la mise à jour. Voir Gérer la qualité des données avec les attentes de pipeline.
Une expectation drop traite une ligne en violation comme si la source ne l'avait jamais produite. La ligne supprimée ne remplace, ne supprime ni ne modifie les clés correspondantes dans la table cible :
- La suppression a lieu avant la déduplication, de sorte que le flux conserve la dernière version valide pour la clé.
- Si chaque ligne entrante pour une clé est supprimée, les lignes existantes de la clé restent inchangées.
- Comme une ligne supprimée ne définit aucun seuil de séquence, une mise à jour valide ultérieure est toujours prise en compte, même si sa séquence est inférieure à celle de la ligne supprimée.
Limitations
Les flux REPLACE USING présentent les limitations suivantes :
- REPLACE USING prend en charge un seul flux par table cible. La combinaison de REPLACE USING avec un autre type de flux sur la même cible n’est pas prise en charge.
- La table cible doit être créée au sein du pipeline.
- La source doit être une source de streaming.
- Vous devez spécifier au moins une colonne clé et une colonne
SEQUENCE BY. Les colonnes clés ne peuvent pas être répétées et le type de chaque colonne clé doit être triable. Les types atomiques, tels que les entiers, les chaînes de caractères et les dates, peuvent être des clés, contrairement àMAPetVARIANT.
Exemples
Les exemples suivants lisent à partir de samples.wanderbricks.booking_updates, une table d'exemple de changements d'état de réservation disponible dans chaque Workspace activé pour Unity Catalog. Chaque réservation apparaît une fois par changement, donc booking_id se répète avec un nouveau booking_update_id. Voir le dataset Wanderbricks.
Exemple 1 : conserver le dernier enregistrement pour chaque clé
Conservez uniquement l'état actuel de chaque réservation. Le flux utilise des clés sur booking_id et des séquences par booking_update_id, de sorte que la mise à jour la plus récente d'une réservation remplace les précédentes. Utilisez plutôt AUTO CDC lorsque votre source est un flux de modifications avec des opérations explicites d'insertion, de mise à jour et de suppression.
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);
from pyspark import pipelines as dp
@dp.table(
name="bookings_current",
replace_using=["booking_id"],
sequence_by="booking_update_id"
)
def bookings_current():
return spark.readStream.table("samples.wanderbricks.booking_updates")
Cet exemple utilise booking_update_id pour la séquence plutôt que le Timestamp updated_at, car plusieurs mises à jour d’une même réservation peuvent partager un même Timestamp. Les lignes qui correspondent à la séquence sont ajoutées plutôt que remplacées, ce qui laisserait plus d’une ligne pour ces réservations.
Exemple 2 : Clé sur plus d’une colonne
Lorsqu’un enregistrement est identifié par une combinaison de colonnes, listez-les toutes dans REPLACE USING. Ici, chaque réservation est identifiée par (property_id, booking_id), de sorte que le flux conserve l’état actuel de chaque réservation par propriété. Si une colonne clé peut être nulle, REPLACE USING fait correspondre null à null au lieu de sauter la ligne.
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);
from pyspark import pipelines as dp
@dp.table(
name="bookings_by_property",
replace_using=["property_id", "booking_id"],
sequence_by="booking_update_id"
)
def bookings_by_property():
return spark.readStream.table("samples.wanderbricks.booking_updates")
Exemple 3 : supprimer les enregistrements non valides avec une attente
Ajoutez une attente pour empêcher les lignes incorrectes d’atteindre la cible. Une ligne supprimée est traitée comme si la source ne l’avait jamais produite : elle ne remplace ni ne supprime la clé correspondante, et le flux revient à la dernière ligne valide pour cette clé. Ce flux ignore les mises à jour qui n’ont pas de total_amount positif.
from pyspark import pipelines as dp
@dp.table(
name="bookings_validated",
replace_using=["booking_id"],
sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
return spark.readStream.table("samples.wanderbricks.booking_updates")