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.
Lorsque plusieurs blocs-notes Fabric, pipelines ou travaux Spark écrivent dans la même table Delta en même temps, Delta Lake utilise le contrôle d’accès concurrentiel optimiste (OCC) pour maintenir la cohérence de la table. Chaque transaction lit un instantané, écrit de nouveaux fichiers, puis valide qu’aucune validation en conflit n’a eu lieu entre elles. Si un conflit est détecté, la transaction échoue en générant une exception plutôt qu’en corrompant les données.
Cet article présente des modèles pratiques pour la gestion des écritures simultanées dans Fabric. Pour obtenir une spécification complète du protocole OCC, consultez le contrôle d’accès concurrentiel Delta Lake (documentation open source).
Niveaux d’isolation
Toutes les tables Delta utilisent le niveau d’isolation Serializable. Serializable est le niveau le plus strict et le seul pris en charge. Il garantit que le résultat des transactions simultanées est identique à un ordre d’exécution séquentiel.
Delta Lake utilise également un niveau SnapshotIsolation interne pour les opérations qui ne modifient pas les données logiques (par OPTIMIZEexemple). SnapshotIsolation ignore la vérification de l’ajout simultané, ce qui permet de poursuivre le compactage sans conflit avec les insertions simultanées. Vous ne configurez pas SnapshotIsolation directement : Delta Lake l’applique automatiquement, le cas échéant.
Avec l’isolation Serializable, un ajout aveugle concurrent (INSERT INTO) peut être en conflit avec un MERGE ou UPDATE qui lit la même partition.
Quelles opérations sont en conflit
Toutes les écritures simultanées n’entrent pas en conflit. Le facteur clé est de savoir si deux opérations touchent les mêmes fichiers sous-jacents.
| Paire concurrente | Conflit? | Pourquoi |
|---|---|---|
Deux INSERT opérations d’ajout |
Non | Chacun ajoute de nouveaux fichiers sans lire les fichiers existants (ajout aveugle). |
INSERT + OPTIMIZE |
Non |
OPTIMIZE est validé à SnapshotIsolation puisqu’il ne modifie pas les données logiques, et ignore donc entièrement la vérification de l’ajout concurrent. Les ajouts créent de nouveaux fichiers qui ne se recoupent pas avec les fichiers en cours de compaction. |
Deux UPDATE, DELETEou MERGE opérations |
Oui, s’ils lisent ou modifient des fichiers qui se chevauchent | Chacun réécrit les fichiers, donc l’instantané du second processus d’écriture n’est plus à jour. |
OPTIMIZE + UPDATE/DELETE/MERGE |
Oui, s’ils touchent les mêmes fichiers |
OPTIMIZE supprime et lit les fichiers (avec dataChange=false). Si une opération de modification de données lit également ces mêmes fichiers, une ConcurrentDeleteReadException est générée. |
Deux OPTIMIZE exécutions |
Oui, s’ils sélectionnent les mêmes fichiers | Les deux tentatives de suppression et de réécriture du même jeu de fichiers déclenchent un ConcurrentDeleteDeleteException. |
INSERT + MERGE/UPDATE/DELETE |
Oui, si l’opération de modification des données lit la même partition | Sous Serializable, un ajout aveugle peut entrer en conflit avec les modifications de données simultanées si l’opération lit une partition dans laquelle l’ajout a écrit. |
Tip
Les pipelines en ajout seul (INSERT INTO, df.write.mode("append")) sont le moyen le plus simple d’éviter entièrement les conflits. Si votre charge de travail peut d’abord ajouter, puis effectuer la réconciliation plus tard, vous éliminez la contention entre écritures.
Isoler les processus d’écriture grâce au partitionnement
La façon la plus courante d’exécuter simultanément DML sur la même table sans conflit consiste à partitionner la table par la colonne qui sépare vos enregistreurs, puis à inclure cette colonne dans chaque condition d’opération. Lorsque chaque enregistreur cible une partition différente, les opérations touchent les jeux de fichiers disjoints et ne sont pas en conflit.
Un scénario typique : plusieurs pipelines traitent chacun des données pour une unité commerciale ou un locataire différent. Partitionnez selon cette dimension et épinglez le MERGE de chaque pipeline à sa partition.
-- Each pipeline targets its own partition, so concurrent runs don't conflict
MERGE INTO events AS target
USING staged AS source
ON target.event_id = source.event_id
AND target.business_unit = 'EMEA'
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
Important
La colonne de partition doit apparaître dans la condition de fusion elle-même, pas seulement dans les données sources. Sans cela, Delta Lake ne peut pas déterminer au moment de la validation que les deux opérations touchaient des jeux de fichiers disjoints, et le vérificateur de conflit traite l’opération comme une lecture de table complète.
Pour plus d’informations sur les stratégies de partitionnement, consultez Partitionnement pour les tables Delta.
Nouvelle tentative de commit intégrée
Delta Lake retente automatiquement une validation lorsqu’elle détecte qu’une autre transaction a été validée en premier. Sur chaque nouvelle tentative, elle lit la validation gagnante, exécute le vérificateur de conflit et, si aucun conflit logique n’existe, réattempte la validation à la prochaine version disponible. Ce processus se répète de manière transparente sans aucune action de votre code.
Un conflit logique (par exemple, deux opérations réécritant le même fichier) ne peuvent pas être résolues automatiquement. La nouvelle tentative génère l’une des exceptions répertoriées dans Exceptions de conflit courantes. Toutefois, de nombreuses collisions de version transitoires — par exemple lorsque deux ajouts sans lecture préalable entrent en concurrence pour le même emplacement de version — sont résolues automatiquement et ne remontent jamais jusqu’à votre application.
Exceptions de conflit courantes
Lorsqu’un conflit est détecté, Delta Lake déclenche une exception spécifique. Comprendre l’exception que vous voyez permet d’identifier la cause racine.
| Exception | Que s’est-il passé |
|---|---|
ConcurrentAppendException |
Un autre enregistreur a ajouté des fichiers dans une partition (ou un jeu de fichiers) que votre opération lisait. Courant lorsqu’une MERGE s’exécute sur une partition recevant également des insertions d’un autre pipeline. Avec l’isolation Serializable, même les ajouts à l’aveugle (simples opérations INSERT) peuvent déclencher cette exception. |
ConcurrentDeleteReadException |
Un autre enregistreur a supprimé ou réécrit un fichier lu par votre opération. Cela se produit généralement lorsque OPTIMIZE compacte des fichiers qu’un processus concurrent UPDATE ou MERGE lisait également, ou lorsque deux opérations de modification de données se chevauchent sur les mêmes lignes. |
ConcurrentDeleteDeleteException |
Les deux opérations ont tenté de supprimer ou de réécrire le même fichier. Souvent provoqué par des OPTIMIZE exécutions qui se chevauchent ou par deux pipelines réécrivant simultanément la même partition. |
ConcurrentWriteException |
Un conflit générique signalé lorsqu’une autre transaction a été validée sur la même version de la table avant que la résolution du conflit ne puisse avoir lieu, par exemple lors d’une mise à niveau d’un système de fichiers vers des validations managées. |
MetadataChangedException |
Le schéma de la table ou ses propriétés ont changé en cours de transaction, par exemple à la suite d’une écriture simultanée ALTER TABLE ou d’une opération d’écriture liée à l’évolution du schéma. |
ConcurrentTransactionException |
Deux requêtes Structured Streaming avec le même emplacement de point de contrôle ont écrit dans la table en même temps. Supprimez les doublons de vos tâches de streaming ou utilisez des chemins de points de contrôle distincts. |
ProtocolChangedException |
Une transaction simultanée a mis à niveau ou rétrogradé le protocole de table pendant que la transaction actuelle tentait également de modifier le protocole. Peut également se produire lorsqu’une fonctionnalité de table est supprimée simultanément. |
Stratégies courantes pour éviter les conflits d’écriture
Activer le compactage automatique
Le compactage automatique s’exécute de manière synchrone dans le cadre des opérations d’écriture. Le compactage synchrone empêche que des tâches de compactage planifiées séparément ne se chevauchent avec des opérations de modification des données et n’entraînent ainsi des exceptions d’écriture concurrente.
Planifier la maintenance en dehors des fenêtres d’écriture
OPTIMIZE et VACUUM peut entrer en conflit avec les opérations de modification de données simultanées. Dans Fabric, planifiez des tâches de notebook ou des activités de pipeline pour le compactage des tables et VACUUM pendant les périodes de faible activité, par exemple après la fin de l’ingestion nocturne, plutôt que pendant celle-ci.
Utiliser des modèles d’ajout + fusion
Pour l’ingestion à forte concurrence, chargez les données brutes avec des écritures en ajout seul dans une table intermédiaire (aucun conflit possible), puis exécutez une seule tâche MERGE pour réconcilier les données dans la table cible. Le modèle sérialise l’opération sujette aux conflits tout en conservant l’ingestion entièrement parallèle.
Ajouter une logique de nouvelle tentative pour les conflits logiques
Le mécanisme intégré de nouvelle tentative de validation gère automatiquement les collisions de version transitoires, mais les conflits logiques — lorsque deux opérations se chevauchent réellement — provoquent une exception. Étant donné que Delta Lake ne produit jamais d’écritures partielles, une transaction ayant échoué est sécurisée pour réessayer au niveau de l’application. Pour les pipelines dans lesquels des conflits logiques ponctuels sont à prévoir, encapsulez l’opération d’écriture dans un mécanisme de nouvelle tentative :
from delta.exceptions import ConcurrentAppendException
import time
# Retry with backoff on transient concurrent write conflicts
max_retries = 3
for attempt in range(max_retries):
try:
spark.sql("MERGE INTO target USING source ON ...")
break
except ConcurrentAppendException:
if attempt < max_retries - 1:
time.sleep(2 ** attempt)
else:
raise
Choisir la stratégie de disposition appropriée
Le clustering liquide et le partitionnement résolvent différents problèmes. Le clustering liquide optimise la disposition des fichiers pour les performances de lecture. Le partitionnement crée des frontières physiques qui empêchent les conflits entre écritures concurrentes. Si votre charge de travail nécessite les deux, partitionnez les données selon la colonne d’isolation en écriture et utilisez Z-Order dans chaque partition afin d’optimiser les performances de lecture.