Remplacement partiel d’un instantané avec des flux REPLACE USING

Important

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 correspondant aux colonnes clés spécifiées et laisse toutes les autres données inchangées.

Une colonne ordonne les mises à jour pour que le résultat soit correct même lorsque les SEQUENCE BY mises à jour arrivent dans le désordre. Pour chaque clé, le numéro de séquence le plus élevé l’emporte, et une ligne avec un numéro de séquence inférieur n’écrase jamais une ligne avec un numéro de séquence plus élevé 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.

Fonctionnement de REPLACE USING

Considérons une table d’événements qui contient des événements de clics et de conversion pour deux régions, séquencés par seq:

region_id type_d’appareil event_type séq.
1 iOS click 1
1 Android conversion 1
2 iOS click 1
2 bureau click 1

Un REPLACE USING (region_id) SEQUENCE BY seq flux reçoit ces mises à jour pour les régions 1 et 3. La Région 2 n’a aucune mise à jour :

region_id type_d’appareil event_type séq.
1 iOS click 2
1 Android conversion 2
1 bureau click 2
3 iOS click 1
3 bureau click 2

La cible devient :

region_id type_d’appareil event_type Suiv Résultat
1 iOS click 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 bureau click 2 Remplacé, car la séquence 2 est supérieure à la séquence 1
2 iOS click 1 Non modifiée, car la clé est absente dans cette mise à jour
2 bureau click 1 Non modifiée, car la clé n’est pas incluse dans cette mise à jour
3 bureau click 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 :

  • REMPLACER EN UTILISANT des flux exécutés sur Databricks Runtime 18.2 et supérieur, sur le calcul classique ou serverless. Databricks recommande Unity Catalog.
  • La source doit être une source de diffusion en continu. REPLACE USING rejette une source qui n’est pas diffusée en continu.
  • Vous devez spécifier au moins une colonne clé et exactement une SEQUENCE BY colonne.

Quand utiliser les flux REPLACE USING

Les pipelines Lakeflow proposent trois flux qui écrasent les lignes existantes. Choisissez en fonction de l’apparence de votre source et de la façon 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 correspondent aux données en entrée, sans modifier toutes les autres données. Il 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 lorsque vous avez besoin d’un historique de dimension de changement lent (SCD) Type 2 . AUTO CDC nécessite également une véritable clé primaire. Consultez les API AUTO CDC : Simplifiez la capture de données modifiées avec des pipelines.
  • Utilisez REPLACE WHERE lorsque votre source est un instantané et que vous souhaitez recalculer et remplacer une plage de la table cible, définie par un prédicat, par exemple les 7 derniers jours, dans le cadre d’un traitement par lots. Il ne nécessite pas de clé primaire. Voir le traitement par lots avec des flux REPLACE WHERE.

Créer un flux de type REMPLACER EN UTILISANT

Définissez les flux REPLACE USING en SQL ou en Python.

SQL

Utilisez la FLOW REPLACE USING clause inline 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);

Vous pouvez également utiliser 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);

Note

BY NAME est requis dans SQL. Il correspond aux colonnes par nom plutôt que par position.

Python

Déclarons la table et le flux 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")

Sinon, 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 lorsque replace_using est défini.

Séquençage et données hors ordre

La SEQUENCE BY colonne 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 à celle déjà stockée pour cette clé, de sorte qu’une ligne tardive ou rejouée plus ancienne que la valeur actuelle est ignorée. Les touches absentes d’une mise à jour restent intactes.

Suivez ces pratiques afin que le remplacement se comporte de manière prévisible :

Pratique Reason
Utilisez une séquence qui augmente strictement par version de clé, comme un horodatage, un numéro de version ou un décalage logarithique. Deux lignes avec 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 indéfini.

Expectations

REMPLACER EN UTILISANT les flux soutiennent les attentes. warn et fail se comportent comme sur d’autres flux : warn ils continuent à violer les lignes, enregistrent la violation et fail arrêtent la mise à jour. Voir Gérer la qualité des données avec les attentes de la chaîne de traitement.

Une drop attente traite une dispute violative comme si la source ne l’avait jamais produite. La ligne supprimée ne remplace, ne supprime pas et ne modifie pas les clés correspondantes dans la table cible :

  • L’élimination a lieu avant la déduplication, de sorte que le flux conserve la dernière version valide associée à la clé.
  • Si chaque ligne entrante d’une clé est supprimée, les lignes existantes de la clé restent intactes.
  • Parce qu’une ligne supprimée ne définit pas de plancher de séquence, une mise à jour valide ultérieure est toujours valide 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 pour chaque 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 dans le pipeline.
  • La source doit être une source de diffusion en continu.
  • Vous devez spécifier au moins une colonne clé et une SEQUENCE BY colonne. 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 et les dates, peuvent être des clés, tandis que MAP et VARIANT ne le peuvent pas.
  • Pour les tables de streaming autonomes, voir Appliquer le remplacement partiel des snapshots à l’aide des flux REPLACE USING pour connaître les différences de syntaxe.

Examples

Les exemples suivants sont tirés de samples.wanderbricks.booking_updates, un tableau d’exemple des changements d’état des réservations disponible dans chaque espace de travail compatible Unity Catalog. Chaque réservation apparaît une fois pour chaque modification, donc booking_id se répète avec un nouveau booking_update_id. Voir le jeu de données Wanderbricks.

Exemple 1 : Conservez le dernier enregistrement pour chaque clé

Ne conservez que l’état actuel de chaque réservation. Le flux est indexé sur booking_id et ordonné selon booking_update_id, de sorte que la mise à jour la plus récente d’une réservation remplace les mises à jour précédentes. Utilisez plutôt AUTO CDC lorsque votre source est un flux de changements avec des opérations explicites d’insertion, de mise à jour et de suppression.

SQL

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

Python

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 est ordonné selon booking_update_id plutôt que selon l’horodatage updated_at, car plusieurs mises à jour d’une même réservation peuvent avoir le même horodatage. Les lignes ayant la même 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 plusieurs colonnes

Lorsqu’un enregistrement est identifié par une combinaison de colonnes, listez-les tous 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 d’ignorer la ligne.

SQL

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

Python

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 à l’aide d’une expectation

Ajoutez une attente pour empêcher les mauvaises lignes d’atteindre la cible. Une ligne perdue 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")