Remplacement partiel de l’instantané par REMPLACER EN UTILISANT des flux

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é, la séquence la plus haute l’emporte, et une ligne de séquence inférieure ne remplace jamais une ligne plus haute déjà 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 le REMPLACEMENT UTILISANT

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 Suiv
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 Suiv
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 Intacte, car la clé n’est pas présente dans cette mise à jour
2 bureau click 1 Intacte, car la clé n’est pas présente 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

REMPLACER EN UTILISANT les flux 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. REMPLACER UTILISER rejette une source non diffusée.
  • Vous devez spécifier au moins une colonne clé et exactement une SEQUENCE BY colonne.

Quand utiliser REMPLACER AVEC les flux

Les pipelines de lakeflow offrent trois écoulements 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 REMPLACER USING lorsque votre source est une série d’instantanés partiels codés par colonne. REMPLACER USING écrase uniquement les données qui correspondent aux données entrantes, laissant toutes les autres données intactes. 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 écraser une plage de la table cible sélectionnée par un prédicat, par exemple les 7 derniers jours, en tant qu’opération batch. Il ne nécessite pas de clé primaire. Voir le traitement par lots avec des flux REPLACE WHERE.

Créer un flux REMPLACER EN UTILISANT

Définissez REMPLACER AVEC des flux 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 Column expression, et est nécessaire chaque fois replace_using que la fonction est définie.

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 :

  • La chute se produit avant la déduplication, donc le flux conserve la dernière version valide de 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

REMPLACER LES flux UTILISANT ont les limitations suivantes :

  • REMPLACER USING prend en charge un seul débit par table cible. Combiner REMPLACER USING avec un autre type de flux sur la même cible n’est pas pris 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 tonalités, tandis que MAP non.VARIANT

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 par changement, donc booking_id elle se répète avec un nouveau booking_update_idfichier . 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 flow s’active booking_id et sequence by booking_update_id, donc la mise à jour la plus récente pour une réservation remplace ses 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 séquence par booking_update_id plutôt que par horodatage updated_at , car plusieurs mises à jour d’une même réservation peuvent partager un horodatage. Les lignes qui s’attachent à la séquence sont ajoutées plutôt qu’remplacées, ce qui laisserait plus d’une rangée pour ces réservations.

Exemple 2 : Clé sur plus d’une colonne

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, REMPLACER EN UTILISANT associe nul·le à nulle plutôt que de sauter 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 dossiers invalides avec une attente

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 laisse tomber les mises à jour qui n’ont pas de positif total_amount.

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