Fusionner des données dans une table Delta Lake à l'aide de l'opération Merge
Vous pouvez insérer/mettre à jour des données à partir d'une table source, d'une vue ou d'un DataFrame dans une table Delta Lake cible en utilisant l'opération SQL MERGE. Delta Lake prend en charge les insertions, les mises à jour et les suppressions dans MERGE, et il prend en charge une syntaxe étendue au-delà des standards SQL pour faciliter les cas d'utilisation avancés.
Supposons que vous ayez une table source nommée people10mupdates ou un chemin source à l'emplacement /tmp/delta/people-10m-updates qui contient de nouvelles données pour une table cible nommée people10m ou un chemin cible à l'emplacement /tmp/delta/people-10m. Certains de ces nouveaux enregistrements peuvent déjà être présents dans les données cibles. To Merge the new data, vous souhaitez mettre à jour les lignes où le id de la personne est déjà présent et insérer les nouvelles lignes où aucun id correspondant n'est présent. Vous pouvez exécuter la query suivante :
- SQL
- Python
- Scala
MERGE INTO people10m
USING people10mupdates
ON people10m.id = people10mupdates.id
WHEN MATCHED THEN
UPDATE SET
id = people10mupdates.id,
firstName = people10mupdates.firstName,
middleName = people10mupdates.middleName,
lastName = people10mupdates.lastName,
gender = people10mupdates.gender,
birthDate = people10mupdates.birthDate,
ssn = people10mupdates.ssn,
salary = people10mupdates.salary
WHEN NOT MATCHED
THEN INSERT (
id,
firstName,
middleName,
lastName,
gender,
birthDate,
ssn,
salary
)
VALUES (
people10mupdates.id,
people10mupdates.firstName,
people10mupdates.middleName,
people10mupdates.lastName,
people10mupdates.gender,
people10mupdates.birthDate,
people10mupdates.ssn,
people10mupdates.salary
)
from delta.tables import *
deltaTablePeople = DeltaTable.forName(spark, "people10m")
deltaTablePeopleUpdates = DeltaTable.forName(spark, "people10mupdates")
dfUpdates = deltaTablePeopleUpdates.toDF()
deltaTablePeople.alias('people') \
.merge(
dfUpdates.alias('updates'),
'people.id = updates.id'
) \
.whenMatchedUpdate(set =
{
"id": "updates.id",
"firstName": "updates.firstName",
"middleName": "updates.middleName",
"lastName": "updates.lastName",
"gender": "updates.gender",
"birthDate": "updates.birthDate",
"ssn": "updates.ssn",
"salary": "updates.salary"
}
) \
.whenNotMatchedInsert(values =
{
"id": "updates.id",
"firstName": "updates.firstName",
"middleName": "updates.middleName",
"lastName": "updates.lastName",
"gender": "updates.gender",
"birthDate": "updates.birthDate",
"ssn": "updates.ssn",
"salary": "updates.salary"
}
) \
.execute()
import io.delta.tables._
import org.apache.spark.sql.functions._
val deltaTablePeople = DeltaTable.forName(spark, "people10m")
val deltaTablePeopleUpdates = DeltaTable.forName(spark, "people10mupdates")
val dfUpdates = deltaTablePeopleUpdates.toDF()
deltaTablePeople
.as("people")
.merge(
dfUpdates.as("updates"),
"people.id = updates.id")
.whenMatched
.updateExpr(
Map(
"id" -> "updates.id",
"firstName" -> "updates.firstName",
"middleName" -> "updates.middleName",
"lastName" -> "updates.lastName",
"gender" -> "updates.gender",
"birthDate" -> "updates.birthDate",
"ssn" -> "updates.ssn",
"salary" -> "updates.salary"
))
.whenNotMatched
.insertExpr(
Map(
"id" -> "updates.id",
"firstName" -> "updates.firstName",
"middleName" -> "updates.middleName",
"lastName" -> "updates.lastName",
"gender" -> "updates.gender",
"birthDate" -> "updates.birthDate",
"ssn" -> "updates.ssn",
"salary" -> "updates.salary"
))
.execute()
Une seule ligne de la table source peut correspondre à une ligne donnée dans la table cible. Dans Databricks Runtime 16.0 et versions ultérieures, MERGE évalue les conditions spécifiées dans les clauses WHEN MATCHED et ON pour identifier les doublons. Dans Databricks Runtime 15,4 LTS et versions antérieures, les opérations MERGE ne prennent en compte que les conditions spécifiées dans la clause ON.
Consultez la documentation de l'API Delta Lake pour les détails de la syntaxe Scala et Python. Pour les détails de la syntaxe SQL, voir MERGE INTO
Modifiez toutes les lignes non correspondantes à l'aide de Merge
Dans Databricks SQL et Databricks Runtime 12.2 LTS et versions ultérieures, vous pouvez utiliser la clause WHEN NOT MATCHED BY SOURCE pour UPDATE ou DELETE des enregistrements dans la table cible qui n'ont pas d'enregistrements correspondants dans la table source. Databricks recommande d'ajouter une clause conditionnelle facultative pour éviter de réécrire entièrement la table cible.
L'exemple de code suivant montre la syntaxe de base pour l'utiliser afin de supprimer, de remplacer la table cible par le contenu de la table source et de supprimer les enregistrements non concordants dans la table cible. Pour un modèle plus évolutif pour les tables où les mises à jour et les suppressions de la source sont limitées dans le temps, consultez Synchroniser une table Delta Lake de manière incrémentielle avec la source.
- Python
- Scala
- SQL
(targetDF
.merge(sourceDF, "source.key = target.key")
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.whenNotMatchedBySourceDelete()
.execute()
)
targetDF
.merge(sourceDF, "source.key = target.key")
.whenMatched()
.updateAll()
.whenNotMatched()
.insertAll()
.whenNotMatchedBySource()
.delete()
.execute()
MERGE INTO target
USING source
ON source.key = target.key
WHEN MATCHED THEN
UPDATE SET *
WHEN NOT MATCHED THEN
INSERT *
WHEN NOT MATCHED BY SOURCE THEN
DELETE
L'exemple suivant ajoute des conditions à la clause WHEN NOT MATCHED BY SOURCE et spécifie les valeurs à mettre à jour dans les lignes cible non correspondantes.
- Python
- Scala
- SQL
(targetDF
.merge(sourceDF, "source.key = target.key")
.whenMatchedUpdate(
set = {"target.lastSeen": "source.timestamp"}
)
.whenNotMatchedInsert(
values = {
"target.key": "source.key",
"target.lastSeen": "source.timestamp",
"target.status": "'active'"
}
)
.whenNotMatchedBySourceUpdate(
condition="target.lastSeen >= (current_date() - INTERVAL '5' DAY)",
set = {"target.status": "'inactive'"}
)
.execute()
)
targetDF
.merge(sourceDF, "source.key = target.key")
.whenMatched()
.updateExpr(Map("target.lastSeen" -> "source.timestamp"))
.whenNotMatched()
.insertExpr(Map(
"target.key" -> "source.key",
"target.lastSeen" -> "source.timestamp",
"target.status" -> "'active'",
)
)
.whenNotMatchedBySource("target.lastSeen >= (current_date() - INTERVAL '5' DAY)")
.updateExpr(Map("target.status" -> "'inactive'"))
.execute()
MERGE INTO target
USING source
ON source.key = target.key
WHEN MATCHED THEN
UPDATE SET target.lastSeen = source.timestamp
WHEN NOT MATCHED THEN
INSERT (key, lastSeen, status) VALUES (source.key, source.timestamp, 'active')
WHEN NOT MATCHED BY SOURCE AND target.lastSeen >= (current_date() - INTERVAL '5' DAY) THEN
UPDATE SET target.status = 'inactive'
Sémantique des opérations de Merge.
Voici une description détaillée de la sémantique de l'opération programmatique merge.
-
Il peut y avoir un nombre quelconque de clauses
whenMatchedetwhenNotMatched. -
whenMatchedLes clauses sont exécutées lorsqu'une ligne source correspond à une ligne de la table cible en fonction de la condition de correspondance. Ces clauses ont la sémantique suivante.-
whenMatchedles clauses peuvent avoir au maximum unupdateet unedeleteaction. L'actionupdatedansmergemet à jour uniquement les colonnes spécifiées (similaire à l'opérationupdate) de la ligne cible correspondante. L'actiondeletesupprime la ligne correspondante. -
Chaque clause
whenMatchedpeut avoir une condition facultative. Si cette condition de clause existe, l'actionupdateoudeleteest exécutée pour toute paire de lignes source-cible correspondante uniquement lorsque la condition de clause est vraie. -
S'il existe plusieurs clauses
whenMatched, alors elles sont évaluées dans l'ordre où elles sont spécifiées. Toutes les clauseswhenMatched, sauf la dernière, doivent avoir des conditions. -
Si aucune des
whenMatchedconditions n'est vraie pour une paire de lignes source et cible qui correspond à la condition de merge, alors la ligne cible reste inchangée. -
Pour mettre à jour toutes les colonnes de la table Delta Lake cible avec les colonnes correspondantes du dataset source, utilisez
whenMatched(...).updateAll(). Ceci est équivalent à :ScalawhenMatched(...).updateExpr(Map("col1" -> "source.col1", "col2" -> "source.col2", ...))pour toutes les colonnes de la table Delta Lake cible. Par conséquent, cette action suppose que la table source possède les mêmes colonnes que celles de la table cible, sinon la query génère une erreur d'analyse.
-
Ce comportement change lorsque l'évolution automatique des schémas est activée. Consultez l'évolution automatique des schémas pour plus de détails.
-
whenNotMatchedles clauses sont exécutées lorsqu'une ligne source ne correspond à aucune ligne cible en fonction de la condition de correspondance. Ces clauses ont la sémantique suivante.-
whenNotMatchedLes clauses ne peuvent avoir que l'actioninsert. La nouvelle ligne est générée en fonction de la colonne spécifiée et des expressions correspondantes. Vous n'avez pas besoin de spécifier toutes les colonnes de la table cible. Pour les colonnes cibles non spécifiées,NULLest inséré. -
Chaque clause
whenNotMatchedpeut avoir une condition facultative. Si la condition de la clause est présente, une ligne source n'est insérée que si cette condition est vraie pour cette ligne. Dans le cas contraire, la colonne source est ignorée. -
S'il existe plusieurs clauses
whenNotMatched, alors elles sont évaluées dans l'ordre où elles sont spécifiées. Toutes les clauseswhenNotMatched, sauf la dernière, doivent avoir des conditions. -
Pour insérer toutes les colonnes de la table Delta Lake cible avec les colonnes correspondantes du dataset source, utilisez
whenNotMatched(...).insertAll(). Ceci est équivalent à :ScalawhenNotMatched(...).insertExpr(Map("col1" -> "source.col1", "col2" -> "source.col2", ...))pour toutes les colonnes de la table Delta Lake cible. Par conséquent, cette action suppose que la table source possède les mêmes colonnes que celles de la table cible, sinon la query génère une erreur d'analyse.
-
Ce comportement change lorsque l'évolution automatique des schémas est activée. Consultez l'évolution automatique des schémas pour plus de détails.
-
whenNotMatchedBySourceles clauses sont exécutées lorsqu’une ligne cible ne correspond à aucune ligne source en fonction de la condition de Merge. Ces clauses ont la sémantique suivante.whenNotMatchedBySourceLes clauses peuvent spécifier des actionsdeleteetupdate.- Chaque clause
whenNotMatchedBySourcepeut avoir une condition facultative. Si la condition de clause est présente, une ligne cible est modifiée uniquement si cette condition est vraie pour cette ligne. Autrement, la ligne cible reste inchangée. - S'il existe plusieurs clauses
whenNotMatchedBySource, alors elles sont évaluées dans l'ordre où elles sont spécifiées. Toutes les clauseswhenNotMatchedBySource, sauf la dernière, doivent avoir des conditions. - Par définition, les clauses
whenNotMatchedBySourcen'ont pas de ligne source à partir de laquelle extraire les valeurs de colonne, et les colonnes source ne peuvent donc pas être référencées. Pour chaque colonne à modifier, vous pouvez soit spécifier un littéral, soit effectuer une action sur la colonne cible, tel queSET target.deleted_count = target.deleted_count + 1.
- Une opération de
mergeMerge peut échouer si plusieurs lignes du dataset source correspondent et que la Merge tente de mettre à jour les mêmes lignes de la table Delta Lake cible. Selon la sémantique SQL de la Merge, une telle opération de mise à jour est ambiguë, car il n'est pas clair quelle ligne source doit être utilisée pour mettre à jour la ligne cible correspondante. Vous pouvez prétraiter la table source pour éliminer la possibilité de correspondances multiples. - Vous pouvez appliquer une opération SQL
MERGEsur une vue SQL uniquement si la vue a été définie commeCREATE VIEW viewName AS SELECT * FROM deltaTable.
Déduplication des données lors de l'écriture dans les tables Delta Lake
Un cas d'utilisation ETL courant consiste à collecter des logs dans une table Delta Lake en les ajoutant à une table. Cependant, souvent, les sources peuvent générer des enregistrements de Logs en double et des étapes de déduplication en aval sont nécessaires pour les gérer. Avec merge, vous pouvez éviter d'insérer les enregistrements en double.
- SQL
- Python
- Scala
- Java
MERGE INTO logs
USING newDedupedLogs
ON logs.uniqueId = newDedupedLogs.uniqueId
WHEN NOT MATCHED
THEN INSERT *
deltaTable.alias("logs").merge(
newDedupedLogs.alias("newDedupedLogs"),
"logs.uniqueId = newDedupedLogs.uniqueId") \
.whenNotMatchedInsertAll() \
.execute()
deltaTable
.as("logs")
.merge(
newDedupedLogs.as("newDedupedLogs"),
"logs.uniqueId = newDedupedLogs.uniqueId")
.whenNotMatched()
.insertAll()
.execute()
deltaTable
.as("logs")
.merge(
newDedupedLogs.as("newDedupedLogs"),
"logs.uniqueId = newDedupedLogs.uniqueId")
.whenNotMatched()
.insertAll()
.execute();
Le dataset contenant les nouveaux Logs doit être dédupliqué en son sein. Selon la sémantique SQL de Merge, il fait correspondre et déduplique les nouvelles données avec les données existantes de la table, mais s'il y a des données en double dans le nouveau dataset, elles sont insérées. Par conséquent, dédupliquez les nouvelles données avant de les fusionner dans la table.
Si vous savez que vous ne pouvez obtenir des enregistrements en double que pendant quelques jours, vous pouvez optimiser davantage votre requête en partitionnant la table par date, puis en spécifiant la plage de dates de la table cible à faire correspondre.
- SQL
- Python
- Scala
- Java
MERGE INTO logs
USING newDedupedLogs
ON logs.uniqueId = newDedupedLogs.uniqueId AND logs.date > current_date() - INTERVAL 7 DAYS
WHEN NOT MATCHED AND newDedupedLogs.date > current_date() - INTERVAL 7 DAYS
THEN INSERT *
deltaTable.alias("logs").merge(
newDedupedLogs.alias("newDedupedLogs"),
"logs.uniqueId = newDedupedLogs.uniqueId AND logs.date > current_date() - INTERVAL 7 DAYS") \
.whenNotMatchedInsertAll("newDedupedLogs.date > current_date() - INTERVAL 7 DAYS") \
.execute()
deltaTable.as("logs").merge(
newDedupedLogs.as("newDedupedLogs"),
"logs.uniqueId = newDedupedLogs.uniqueId AND logs.date > current_date() - INTERVAL 7 DAYS")
.whenNotMatched("newDedupedLogs.date > current_date() - INTERVAL 7 DAYS")
.insertAll()
.execute()
deltaTable.as("logs").merge(
newDedupedLogs.as("newDedupedLogs"),
"logs.uniqueId = newDedupedLogs.uniqueId AND logs.date > current_date() - INTERVAL 7 DAYS")
.whenNotMatched("newDedupedLogs.date > current_date() - INTERVAL 7 DAYS")
.insertAll()
.execute();
C'est plus efficace que la commande précédente, car elle recherche les doublons uniquement dans les 7 derniers jours des Logs, et non dans la table entière. De plus, vous pouvez utiliser cette Merge en insertion uniquement avec Structured Streaming pour effectuer une déduplication continue des Logs.
- Dans une query de streaming, vous pouvez utiliser l'opération de Merge dans
foreachBatchpour écrire en continu des données de streaming dans une table Delta Lake avec déduplication. Consultez l'exemple de streaming suivant pour plus d'informations surforeachBatch. - Dans une autre query de streaming, vous pouvez lire en continu des données dédupliquées de cette table Delta Lake. Ceci est possible car une Merge en mode insertion seule n'ajoute que de nouvelles données à la table Delta Lake.
Données à évolution lente (SCD) et capture des données de modification (CDC) avec Delta Lake
Les LakeFlow Pipelines prennent en charge nativement le suivi et l'application des SCD de Type 1 et de Type 2. Utilisez AUTO CDC ... INTO avec les LakeFlow Pipelines pour vous assurer que les enregistrements désordonnés sont traités correctement lors du traitement des flux CDC. Consultez Les AUTO CDC APIs : simplifiez la capture des données modifiées avec des pipeline.
Synchroniser de manière incrémentielle la table Delta Lake avec la source
Dans Databricks SQL et Databricks Runtime 12.2 LTS et versions ultérieures, vous pouvez utiliser WHEN NOT MATCHED BY SOURCE pour créer des conditions arbitraires afin de supprimer et de remplacer atomiquement une partie d'une table. Cela peut être particulièrement utile lorsque vous avez une table source où les enregistrements peuvent changer ou être supprimés pendant plusieurs jours après la saisie initiale des données, mais finissent par atteindre un état final.
La query suivante montre comment utiliser ce modèle pour sélectionner les enregistrements des 5 derniers jours de la source, mettre à jour les enregistrements correspondants dans la cible, insérer de nouveaux enregistrements de la source vers la cible et supprimer tous les enregistrements non correspondants des 5 derniers jours dans la cible.
MERGE INTO target AS t
USING (SELECT * FROM source WHERE created_at >= (current_date() - INTERVAL '5' DAY)) AS s
ON t.key = s.key
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
WHEN NOT MATCHED BY SOURCE AND created_at >= (current_date() - INTERVAL '5' DAY) THEN DELETE
En fournissant le même filtre booléen sur les tables source et cible, vous êtes capable de propager dynamiquement les modifications de votre source vers les tables cibles, y compris les suppressions.
Bien que ce modèle puisse être utilisé sans aucune clause conditionnelle, cela entraînerait la réécriture complète de la table cible, ce qui peut être coûteux.