Aller au contenu principal

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

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
(targetDF
.merge(sourceDF, "source.key = target.key")
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.whenNotMatchedBySourceDelete()
.execute()
)

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

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 whenMatched et whenNotMatched.

  • whenMatched Les 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.

    • whenMatched les clauses peuvent avoir au maximum un update et une delete action. L'action update dans merge met à jour uniquement les colonnes spécifiées (similaire à l'opération update) de la ligne cible correspondante. L'action delete supprime la ligne correspondante.

    • Chaque clause whenMatched peut avoir une condition facultative. Si cette condition de clause existe, l'action update ou delete est 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 clauses whenMatched, sauf la dernière, doivent avoir des conditions.

    • Si aucune des whenMatched conditions 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 à :

      Scala
      whenMatched(...).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.

remarque

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.

  • whenNotMatched les 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.

    • whenNotMatched Les clauses ne peuvent avoir que l'action insert. 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, NULL est inséré.

    • Chaque clause whenNotMatched peut 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 clauses whenNotMatched, 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 à :

      Scala
      whenNotMatched(...).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.

remarque

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.

  • whenNotMatchedBySource les 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.

    • whenNotMatchedBySource Les clauses peuvent spécifier des actions delete et update.
    • Chaque clause whenNotMatchedBySource peut 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 clauses whenNotMatchedBySource, sauf la dernière, doivent avoir des conditions.
    • Par définition, les clauses whenNotMatchedBySource n'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 que SET target.deleted_count = target.deleted_count + 1.
important
  • Une opération de merge Merge 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 MERGE sur une vue SQL uniquement si la vue a été définie comme CREATE 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
MERGE INTO logs
USING newDedupedLogs
ON logs.uniqueId = newDedupedLogs.uniqueId
WHEN NOT MATCHED
THEN INSERT *
remarque

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
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 *

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 foreachBatch pour é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 sur foreachBatch.
  • 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.

SQL
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.

remarque

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.