Aller au contenu principal

Capture des modifications de données et instantanés

Les data engineers doivent souvent répliquer des données provenant de sources en amont de Databricks, telles que des bases de données relationnelles (Oracle, Postgres, SQL Server), dans Databricks pour l'analytique, le reporting et le Machine Learning. À mesure que les systèmes opérationnels changent, les tables analytiques doivent rester synchronisées avec ces changements.

Certaines équipes ont besoin de refléter l'état actuel de leurs bases de données opérationnelles pour le reporting et l'analytique. D'autres doivent conserver un historique complet des changements pour la possibilité d'audit, les exigences réglementaires ou l'analytique client.

La capture de données modifiées (CDC) traite une base de données comme un ensemble de modifications, plutôt que comme une base de données statique complète. Le diagramme suivant montre que lorsqu'une ligne d'une table source contenant des données d'employés est mise à jour, elle génère un nouvel ensemble de lignes dans un flux CDC qui contient uniquement les modifications. Chaque ligne du flux CDC contient généralement des métadonnées supplémentaires, y compris l'Opérations telle que UPDATE et une colonne qui peut être utilisée pour ordonner de manière déterministe chaque ligne dans le flux CDC afin que vous puissiez gérer les mises à jour désordonnées. Par exemple, la colonne sequenceNum du diagramme suivant détermine l'ordre des lignes dans le flux CDC :

Aperçu de la capture de données modifiées.

La CDC vous permet de visualiser simplement les changements apportés aux données pour des transactions plus simples lors de la mise à jour de la base de données dans un système en aval. Cela vous permet également de consulter l'historique de la base de données, si cela est une exigence.

Le défi est que les systèmes sources fournissent des données dans différents formats. Certains émettent des flux de changements qui capturent les changements individuels (insertions, mises à jour, suppressions). D'autres ne fournissent que des instantanés périodiques de la table entière. Chaque format nécessite des approches de traitement différentes pour maintenir l'exactitude et l'actualisation des tables en aval.

Historiquement, les équipes se sont appuyées sur une logique MERGE INTO personnalisée pour appliquer ces modifications, que ce soit à partir de flux de modifications ou en comparant des instantanés. Cette approche est complexe et source d’erreurs, nécessitant des tables de staging, des fonctions de fenêtre et des hypothèses de séquençage difficiles à comprendre et à maintenir à mesure que les pipelines évoluent.

Bénéfices de la CDC

La capture de changement de données offre plusieurs avantages dans vos charges de travail.

  • Les données modifiées sont généralement plus petites que l'ensemble du dataset, et les modifications peuvent être traitées par les requêtes en aval sous forme de mises à jour incrémentielles des données.
  • Les données de modification peuvent être stockées de manière à vous permettre de reconstituer les enregistrements tels qu'ils étaient à un moment précis, vous offrant ainsi un historique complet pour l’audit, le reporting ponctuel ou l’analyse des tendances.
  • Les données de modification permettent des clés de substitution stables au fil du temps.

Comment les modifications sont appliquées : État actuel ou historique complet des modifications

Les dimensions à évolution lente (SCD) définissent la manière dont les changements en amont sont appliqués et modélisés après leur arrivée dans les tables d'analyse. Les organisations peuvent utiliser différentes approches en fonction de leurs besoins en données. Le SCD de type 1 vous permet d'enregistrer uniquement l'état actuel du dataset. Le SCD de type 2 enregistre l'historique complet des modifications apportées au dataset. Cette section les décrit plus en détail.

SCD Type 1 : État actuel uniquement

Le SCD de type 1 écrase les anciennes données avec les nouvelles données chaque fois que des modifications se produisent, ne conservant que la dernière version de chaque enregistrement. L'historique n'est pas conservé.

Utilisez SCD Type 1 lorsque :

  • Vous n'avez besoin que de l'état actuel des données.
  • Vous souhaitez que les vues matérialisées en aval refresh de manière incrémentielle plutôt qu'elles ne soient entièrement recalculées.
  • Vous avez besoin de clés de substitution stables pour les jointures.

Seule la dernière version des données est disponible dans SCD1. C'est une approche simple et directe que l'on peut considérer comme ne stockant que la table finale. Si un enregistrement passe de Owner à Manager,, seul Manager reste dans la table :

Aperçu de la capture des données de changement SCD de type 1.

SCD de type 2 : suivi historique

Le SCD de type 2 maintient un enregistrement historique complet en créant plusieurs versions de données au fil du temps, chacune avec un Timestamp et des métadonnées. Les colonnes __START_AT et __END_AT définissent la période de validité pour chaque version d'un enregistrement. Les enregistrements actifs ont __END_AT = NULL. Vous pouvez afficher l'état du dataset tel qu'il était à tout moment.

Utilisez le SCD de type 2 lorsque :

  • L'auditabilité ou les exigences réglementaires exigent un suivi historique.
  • L'analytique client nécessite de comprendre comment les entités ont évolué au fil du temps.
  • La logique métier nécessite un reporting ponctuel.
  • Vous devez analyser les tendances ou comparer des états historiques.

Le traitement SCD de type 2 maintient un historique des changements de données. Par exemple, si un enregistrement a actuellement le champ rôle défini sur Manager, vous pouvez également constater que le rôle était précédemment défini sur Owner. Dans l'image suivante, c'est exactement ce qui est arrivé à l'enregistrement pour Chris. Vous pouvez identifier l'enregistrement actuel, car il a une valeur null pour le champ end_at :

Aperçu de la capture des données de changement SCD de type 2.

Qu'est-ce qu'un flux CDC ?

La capture de données modifiées (CDC) est un modèle d'intégration de données qui capture les modifications apportées aux données dans un système source : insertions, mises à jour et suppressions. Plutôt que de traiter des dataset entiers, CDC génère des flux contenant uniquement les enregistrements modifiés.

Par exemple, si vous avez une table d’employés dans Oracle avec 50 lignes, et que le titre de poste d’un employé change, le flux CDC contient un seul enregistrement UPDATE pour cet employé. Cela permet à Databricks de traiter uniquement les enregistrements modifiés plutôt que de lire l'intégralité de la table source à chaque exécution.

Chaque enregistrement CDC de la base de données source inclut :

  • Le type d'opération (INSERT, UPDATE, DELETE)
  • Les valeurs de données pour l'enregistrement
  • Un numéro de séquence ou un Timestamp pour un ordonnancement déterministe

Le numéro de séquence garantit que les arrivées tardives ou désordonnées sont appliquées correctement. Les bases de données transactionnelles telles que SQL Server, MySQL et Oracle génèrent des flux CDC en mode natif. Les tables Delta génèrent également leur propre flux CDC, connu sous le nom de flux de données de modification (CDF), facilitant ainsi le traitement des modifications provenant également des sources Delta.

Qu’est-ce qu’un snapshot ?

Un instantané représente l'état complet d'une table à un moment précis. Contrairement aux flux CDC qui ne capturent que les modifications, les instantanés contiennent chaque ligne de la table source.

Les équipes n'activent pas toujours les flux CDC sur les bases de données opérationnelles pour diverses raisons :

  • Coût (la CDC peut augmenter la charge sur les bases de données de production)
  • Problèmes de performances sur la base de données source
  • Systèmes hérités qui ne prennent pas en charge la CDC
  • Contraintes organisationnelles (les équipes gérant l'ingestion ne sont pas propriétaires des bases de données en amont)

Lorsqu'un flux de modifications n'est pas disponible, l'ingestion basée sur des instantanés est la seule option. Les instantanés peuvent provenir de :

  • Exportations périodiques à partir de bases de données relationnelles (Oracle, Postgres, SQL Server)
  • Dumps de fichiers de stockage cloud provenant de systèmes en amont.
  • tables Delta (chaque version de table est effectivement un instantané)
  • OpenSharing depuis des tenants en amont

Étant donné que les instantanés ne capturent pas les modifications au niveau des enregistrements, l'identification de ce qui a changé nécessite de comparer les enregistrements entre les instantanés pour en déduire les insertions, les mises à jour et les suppressions.

Traiter automatiquement les flux CDC

Databricks simplifie le traitement CDC via l'API AUTO CDC au sein des Lakeflow pipelines. Cette API est conçue pour traiter les modifications provenant des flux CDC sur les bases de données sources ou les tables Delta avec le Change Data Feed activé.

Pour des exemples de code SQL et Python, consultez les exemples AUTO CDC.

Utilisez AUTO CDC lorsque l'une de ces conditions est vraie :

  • Votre système source génère un flux de données de changement (CDF).
  • Vous lisez à partir d'une table Delta avec Change Data Feed activé.
  • Vous disposez d'un flux CDC provenant d'une base de données relationnelle (via des outils tels que Debezium ou Oracle GoldenGate)

AUTO CDC gère automatiquement les enregistrements hors séquence en traitant les événements dans l'ordre défini par la colonne de séquencement. La colonne de séquencement doit être une représentation croissante monotone de l'ordre correct des événements, avec une mise à jour distincte par clé à chaque valeur de séquencement. Les valeurs de séquençage NULL ne sont pas prises en charge. Pour le SCD de type 2, le pipeline propage les valeurs de séquençage aux colonnes __START_AT et __END_AT de la table cible.

**Hydratation initiale** : lors de la réplication d'une table de base de données opérationnelle existante dans Databricks, vous devez d'abord charger toutes les données historiques avant de traiter les modifications en cours. AUTO CDC prend en charge cela via les *flux ponctuels*, un mode qui traite toutes les données disponibles une seule fois, puis s'arrête. Une fois le chargement initial terminé, utilisez un flux en mode Trigger ou continu pour le traitement CDC en cours. Ceci assure une logique cohérente pour les chargements en masse et incrémentiels.

Traiter automatiquement les instantanés

Lorsque les flux CDC ne sont pas disponibles, Databricks fournit l'API AUTO CDC FROM SNAPSHOT. Cette API est conçue pour l'ingestion basée sur des snapshots ; elle compare des snapshots consécutifs , génère un flux de modifications synthétique et applique une logique SCD de type 1 ou de type 2 dans la table cible. La table cible peut fournir un flux CDC (appelé *flux de données de modification* (CDF) dans les tables Delta) de type SCD 1 ou de type 2 pour les requêtes en aval.

Pour des exemples de code Python, consultez les exemples AUTO CDC FROM SNAPSHOT.

AUTO CDC FROM SNAPSHOT est pris en charge uniquement dans l'interface de pipeline Python. Les instantanés doivent être traités par ordre croissant de version ; si un instantané hors service est détecté, il est ignoré. Le traitement en aval, comme une vue matérialisée qui query la sortie d'un dataset AUTO CDC FROM SNAPSHOT, bénéficie des avantages du CDC, tels que la capacité d'incrémentalisation et des clés de substitution stables.

remarque

AUTO CDC FROM SNAPSHOT n'est pas seulement pour les chargements initiaux. Il est conçu pour un traitement continu lorsque les instantanés sont votre seul format disponible. Chaque fois qu'un nouvel instantané arrive, l'API le compare à l'instantané précédent pour en déduire les modifications et un flux de données de modification.

Utilisez AUTO CDC FROM SNAPSHOT lorsque :

  • La CDC n'est pas activée sur la base de données source.
  • Vous avez uniquement accès à des instantanés périodiques (dumps de table complets)
  • Vous souhaitez les avantages de la CDC pour le traitement incrémental, ou pour avoir un historique complet des modifications.

AUTO CDC FROM SNAPSHOT gère automatiquement les éléments suivants :

  1. Compare les snapshots consécutifs pour identifier les enregistrements insérés, mis à jour et supprimés.
  2. Génère un flux de modifications synthétique basé sur les différences entre les instantanés.
  3. Applique la même logique SCD que AUTO CDC pour calculer le SCD de type 1 ou de type 2.
remarque

AUTO CDC FROM SNAPSHOT ne connaît que les changements d'un instantané à l'autre et n'obtient pas les changements intermédiaires. Par exemple, si vous obtenez des instantanés quotidiens et qu'un utilisateur change d'adresse deux fois en une journée (de A à B, puis de B à C), votre flux de modifications peut passer directement de A à C, car vous n'avez reçu que des instantanés pour ces moments.

Modèles de traitement Snapshot

AUTO CDC FROM SNAPSHOT prend en charge deux modèles pour déterminer les versions d'instantanés.

Traitement des instantanés utilisant le temps d'ingestion du pipeline

L'instantané est lu au moment de l'exécution du pipeline, et l'heure d'ingestion est utilisée comme version de l'instantané. Un nouvel instantané est ingéré à chaque mise à jour du pipeline. Lorsqu'un pipeline s'exécute en mode continu, plusieurs instantanés sont ingérés en fonction du paramètre d'intervalle du Trigger pour le flux.

Utilisez ce modèle lorsque les instantanés arrivent régulièrement et dans l'ordre et que vous pouvez vous fier au Timestamp d'exécution du pipeline pour la gestion des versions.

Traitement d'instantanés à l'aide de fonctions de version

Vous fournissez une fonction qui spécifie quelle version d'instantané traiter au moment de l'exécution du pipeline. La fonction renvoie un tuple : (DataFrame, version_number). L'API traite les instantanés dans l'ordre défini par les numéros de version. Si un instantané désordonné est détecté, il est ignoré.

Utilisez ce modèle lorsque :

  • Plusieurs snapshots peuvent arriver en même temps et nécessitent un traitement séquentiel.
  • Les instantanés peuvent arriver dans le désordre.
  • Vous avez besoin d'un contrôle explicite sur l'ordonnancement des instantanés.

Fonctionnalités CDC supplémentaires

Opérations de modification sur les cibles AUTO CDC

Contrairement aux tables de streaming standard, les tables Unity Catalog qui sont des cibles AUTO CDC prennent en charge les instructions INSERT, UPDATE, DELETE et MERGE même lorsque le pipeline est en cours d'exécution. Pour plus de détails et de limitations, consultez Ajouter, modifier ou supprimer des données dans une table de streaming cible.

Lecture des flux de données de modification à partir des cibles AUTO CDC

AUTO CDC Les tables streaming cibles peuvent émettre leur propre flux de données de modification (CDF), permettant aux pipelines en aval de consommer les modifications de la sortie AUTO CDC. Pour plus de détails, consultez Lire un flux de données de modification à partir d'une table cible AUTO CDC.

Métriques et monitoring

AUTO CDC capture automatiquement les métriques num_upserted_rows et num_deleted_rows à chaque exécution de pipeline. Pour plus de détails, consultez les sujets CDC AUTO avancés.

Suivi des sous-ensembles de colonnes pour le SCD de Type 2

Par défaut, le SCD de type 2 crée une nouvelle version chaque fois que la valeur d'une colonne change. AUTO CDC vous permet de spécifier les colonnes à suivre pour l'historique, afin que les modifications apportées aux colonnes non suivies mettent à jour la version actuelle sur place plutôt que de créer un nouvel enregistrement historique. Cela réduit les coûts de stockage et la complexité des query tout en préservant l'historique des attributs critiques. Pour un exemple, voir Suivre un sous-ensemble de colonnes avec le SCD de type 2.

Recommandations

Utilisez la capture de données modifiées (CDC) lorsque vous souhaitez ne travailler qu'avec les modifications de vos données, par exemple, pour permettre la mise à jour incrémentielle des vues matérialisées en aval. Utilisez également la CDC lorsque vous souhaitez conserver un historique des modifications apportées à vos données, par exemple, pour savoir qui occupait quel rôle à un moment précis.

Utilisez les API AUTO CDC lorsque vous devez répliquer les données en amont dans Databricks et les maintenir synchronisées avec les modifications de la source. La bonne API dépend de la manière dont votre système source expose les changements :

  • Utilisez AUTO CDC lorsque votre source émet un Stream de modifications, par exemple une base de données relationnelle avec CDC activé (via des outils tels que Debezium ou Oracle GoldenGate), une table Delta avec Change Data Feed activé, ou toute source qui produit un Stream d'insertions, de mises à jour et de suppressions avec une colonne de séquencement.
  • Utilisez AUTO CDC FROM SNAPSHOT lorsque votre source ne prend pas en charge la CDC et ne fournit que des vidages complets périodiques de table. Cette API infère les modifications en comparant des instantanés consécutifs et génère un flux de modifications synthétique, de sorte que vous obtenez les mêmes avantages de traitement SCD même sans un flux CDC natif.

Dans les deux cas, choisissez le type de SCD 1 si vous n'avez besoin que de l'état actuel de chaque enregistrement, ou le type de SCD 2 si vous devez conserver un historique complet des modifications pour l'audit, les rapports à un instant T ou l'analyse des tendances.

Ressources supplémentaires