Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
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 BYcolonne.
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 BYcolonne. 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 queMAPetVARIANTne 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")