Les API AUTO CDC : Simplifier la capture des changements de données avec des pipelines
Les pipelines Lakeflow simplifient la capture des changements de données (CDC) avec les APIs AUTO CDC et AUTO CDC FROM SNAPSHOT. Ces APIs automatisent la complexité du calcul des dimensions à évolution lente (SCD) de Type 1 et de Type 2 à partir d'un flux CDC ou d'instantanés de base de données. L'API AUTO CDC prend également en charge le suivi bitemporel, qui enregistre les changements sur deux dimensions temporelles (Beta). Pour en savoir plus sur les SCD de type 1 et de type 2, consultez Capture des données modifiées et instantanés. Pour en savoir plus sur le suivi bitemporel, consultez Bitemporal AUTO CDC.
Les AUTO CDC APIs remplacent les APPLY CHANGES APIs et ont la même syntaxe. Les APIs APPLY CHANGES sont toujours disponibles, mais Databricks recommande d'utiliser les APIs AUTO CDC à leur place.
L'API que vous utilisez dépend de la source de vos données de modification :
AUTO CDC: Utilisez ceci lorsque la base de données source a un flux CDC activé.AUTO CDCtraite les modifications d'un flux de données de modification (CDF). Il est pris en charge à la fois dans les interfaces SQL et Python du pipeline.AUTO CDC FROM SNAPSHOT: Utilisez ceci lorsque la CDC n'est pas activée sur la base de données source et que seuls des instantanés sont disponibles. Cette API compare les instantanés pour déterminer les changements, puis les traite. Il est pris en charge uniquement dans l'interface Python.
Les deux APIs prennent en charge la mise à jour des tables à l'aide de SCD de type 1 et de type 2 :
- Utilisez le SCD de type 1 pour mettre à jour les enregistrements directement. L'historique n'est pas conservé pour les enregistrements mis à jour.
- Utilisez le SCD de type 2 pour conserver un historique des enregistrements, soit sur toutes les mises à jour, soit sur les mises à jour d'un ensemble de colonnes spécifié.
Pour AUTO CDC uniquement, vous pouvez également utiliser le stockage bitemporel, qui étend l'historique des SCD de Type 2 pour suivre les modifications sur deux dimensions temporelles : le temps métier et le temps système. Bitemporal est en bêta. Voir Bitemporal AUTO CDC.
AUTO CDC prend également en charge les mises à jour partielles, où un enregistrement de modification met à jour seulement un sous-ensemble de colonnes. Consultez Appliquer les mises à jour partielles.
Les APIs AUTO CDC ne sont pas prises en charge par Apache Spark Declarative Pipelines.
Pour la syntaxe et d’autres références, consultez AUTO CDC INTO (pipelines), create_auto_cdc_flow, et create_auto_cdc_from_snapshot_flow.
Cette page décrit comment mettre à jour les tables dans vos pipelines en fonction des modifications des données source. Pour savoir comment enregistrer et interroger les informations de modification au niveau des lignes pour les tables Delta, consultez Utiliser le flux de données de modification sur Databricks.
Exigences
Pour utiliser les API CDC, votre pipeline doit être configuré pour utiliser les Serverless LakeFlow pipelines ou Pro les Advanced éditions ou des LakeFlow pipelines.
Fonctionnement de l'AUTO CDC
Pour effectuer le traitement CDC avec AUTO CDC, créez une table de streaming, puis utilisez l'instruction AUTO CDC ... INTO en SQL ou la fonction create_auto_cdc_flow() en Python pour spécifier la source, les clés et le séquencement du flux de modification. Pour une explication du fonctionnement du séquencement et de la logique SCD, consultez Change data capture and snapshots. Consultez les exemples AUTO CDC.
Pour une hydratation initiale à partir d'une source avec un flux de modifications, utilisez AUTO CDC avec un flux once et continuez ensuite à traiter le flux de modifications. Consultez Répliquer une table RDBMS externe à l'aide d'AUTO CDC.
Pour plus de détails sur la syntaxe, consultez AUTO CDC INTO (pipelines) ou create_auto_cdc_flow.
Comment AUTO CDC FROM SNAPSHOT fonctionne
AUTO CDC FROM SNAPSHOT détermine les changements dans les données source en comparant des instantanés ordonnés. Il est pris en charge uniquement dans l'interface de pipeline Python. Vous pouvez lire les instantanés directement à partir d'une table Delta, de fichiers de stockage cloud ou de JDBC.
Pour effectuer le traitement CDC avec AUTO CDC FROM SNAPSHOT, créez une table de streaming, puis utilisez la fonction create_auto_cdc_from_snapshot_flow() pour spécifier l'instantané, les clés et d'autres arguments. Pour plus de détails sur les deux modèles d'ingestion et quand utiliser chacun, consultez Modèles de traitement d'instantanés. Voir les exemples AUTO CDC FROM SNAPSHOT.
Pour plus de détails sur la syntaxe, consultez create_auto_cdc_from_snapshot_flow.
Utiliser plusieurs colonnes pour le séquençage
Pour séquencer par plusieurs colonnes (par exemple, un Timestamp et un ID pour départager), utilisez un STRUCT pour les combiner. L'API classe par le premier champ, et en cas d'égalité, prend en compte le second champ, et ainsi de suite.
- SQL
- Python
SEQUENCE BY STRUCT(timestamp_col, id_col)
sequence_by = struct("timestamp_col", "id_col")
Exemples AUTO CDC
Les exemples suivants illustrent le traitement SCD de type 1 et de type 2 à l'aide d'une source de flux de données de modifications. Les données d'échantillon créent de nouveaux enregistrements d'utilisateur, suppriment un enregistrement d'utilisateur et mettent à jour les enregistrements d'utilisateur. Dans l'exemple de SCD de type 1, les UPDATE dernières opérations arrivent en retard et sont supprimées de la table cible, ce qui démontre la gestion des événements désordonnés.
Voici les enregistrements d'entrée utilisés dans ces exemples. Ces données sont créées en exécutant la query dans la section Créer un échantillon de données.
ID utilisateur | Nom | ville | Opérations | sequenceNum |
|---|---|---|---|---|
124 | Raul | Oaxaca | INSÉRER | 1 |
123 | Isabel | Monterrey | INSÉRER | 1 |
125 | Mercedes | Tijuana | INSÉRER | 2 |
126 | Lily | Cancun | INSÉRER | 2 |
123 | nul | nul | Supprimer | 6 |
125 | Mercedes | Guadalajara | Mettre à jour | 6 |
125 | Mercedes | Mexicali | Mettre à jour | 5 |
123 | Isabel | Chihuahua | Mettre à jour | 5 |
Si vous décommentez la dernière ligne dans la query de génération de données d'exemple, cela insère l'enregistrement suivant qui spécifie de tronquer la table (effacer la table) à sequenceNum=3:
ID utilisateur | Nom | ville | Opérations | sequenceNum |
|---|---|---|---|---|
nul | nul | nul | TRONQUER | 3 |
Tous les exemples suivants incluent des options pour spécifier les opérations DELETE et TRUNCATE, mais chacune est facultative.
Créer un échantillon de données
Exécutez les instructions suivantes pour créer un jeu de données d'exemple. Ce code n'est pas destiné à être exécuté dans le cadre d'une définition de pipeline. Exécutez-le à partir du dossier d'exploration de votre pipeline, plutôt qu'à partir du dossier de transformations.
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.users_cdf
AS SELECT
col1 AS userId,
col2 AS name,
col3 AS city,
col4 AS operation,
col5 AS sequenceNum
FROM (
VALUES
-- Initial load.
(124, "Raul", "Oaxaca", "INSERT", 1),
(123, "Isabel", "Monterrey", "INSERT", 1),
-- New users.
(125, "Mercedes", "Tijuana", "INSERT", 2),
(126, "Lily", "Cancun", "INSERT", 2),
-- Isabel is removed from the system and Mercedes moved to Guadalajara.
(123, null, null, "DELETE", 6),
(125, "Mercedes", "Guadalajara", "UPDATE", 6),
-- This batch of updates arrived out of order. The batch at sequenceNum 6 is the final state.
(125, "Mercedes", "Mexicali", "UPDATE", 5),
(123, "Isabel", "Chihuahua", "UPDATE", 5)
-- Uncomment to test TRUNCATE.
-- ,(null, null, null, "TRUNCATE", 3)
);
Traiter les mises à jour SCD de type 1
Le SCD de type 1 conserve uniquement la dernière version de chaque enregistrement. L'exemple suivant lit à partir du flux de données de modification créé ci-dessus et applique les modifications à une cible de table de streaming. Que sont les pipelines ? pour exécuter ce code.
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_current")
dp.create_auto_cdc_flow(
target = "users_current",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
apply_as_truncates = expr("operation = 'TRUNCATE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = 1
)
CREATE OR REFRESH STREAMING TABLE users_current;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_current
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
APPLY AS TRUNCATE WHEN
operation = "TRUNCATE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 1;
Après avoir exécuté l'exemple de SCD de type 1, la table cible contient les enregistrements suivants :
ID utilisateur | Nom | ville |
|---|---|---|
124 | Raul | Oaxaca |
125 | Mercedes | Guadalajara |
126 | Lily | Cancun |
L'utilisateur 123 (Isabel) a été supprimé et n'apparaît pas. L'utilisateur 125 (Mercedes) ne montre que la dernière ville (Guadalajara) car le SCD de type 1 écrase les valeurs précédentes. L'élément UPDATE précédent à sequenceNum=5 a été abandonné car une mise à jour ultérieure à sequenceNum=6 est arrivée.
Après avoir exécuté l'exemple avec l'enregistrement TRUNCATE non commenté, la table est effacée à sequenceNum=3. Cela signifie que les enregistrements 124 et 126 ne sont pas dans la table, et que la table cible finale contient uniquement l'enregistrement suivant :
ID utilisateur | Nom | ville |
|---|---|---|
125 | Mercedes | Guadalajara |
Traiter les mises à jour SCD de type 2
Le SCD de type 2 préserve un historique complet des modifications en créant de nouvelles lignes pour chaque version d'un enregistrement, avec les colonnes __START_AT et __END_AT indiquant quand chaque version était active.
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2"
)
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 2;
Après avoir exécuté l'exemple de SCD de type 2, la table cible contient les enregistrements suivants :
ID utilisateur | Nom | ville | start | __END_AT |
|---|---|---|---|---|
123 | Isabel | Monterrey | 1 | 5 |
123 | Isabel | Chihuahua | 5 | 6 |
124 | Raul | Oaxaca | 1 | nul |
125 | Mercedes | Tijuana | 2 | 5 |
125 | Mercedes | Mexicali | 5 | 6 |
125 | Mercedes | Guadalajara | 6 | nul |
126 | Lily | Cancun | 2 | nul |
La table conserve l'historique complet. L'utilisateur 123 a deux versions (terminé à la séquence 6 lors de la suppression). L'utilisateur 125 dispose de trois versions présentant des changements de ville. Les enregistrements avec __END_AT = null sont actuellement actifs.
Suivre un sous-ensemble de colonnes avec 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. Vous pouvez spécifier un sous-ensemble de colonnes à suivre, de sorte que les modifications apportées aux autres colonnes mettent à jour la version actuelle sur place plutôt que de générer un nouvel enregistrement d’historique.
L'exemple suivant exclut la colonne city du suivi de l'historique :
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2",
track_history_except_column_list = ["city"]
)
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 2
TRACK HISTORY ON * EXCEPT
(city)
Étant donné que les modifications de city ne sont pas suivies, les mises à jour de ville écrasent la ligne actuelle au lieu de créer une nouvelle version. La table cible contient les enregistrements suivants :
ID utilisateur | Nom | ville | start | __END_AT |
|---|---|---|---|---|
123 | Isabel | Chihuahua | 1 | 6 |
124 | Raul | Oaxaca | 1 | nul |
125 | Mercedes | Guadalajara | 2 | nul |
126 | Lily | Cancun | 2 | nul |
Exemples d'AUTO CDC FROM SNAPSHOT
Les sections suivantes fournissent des exemples d'utilisation de AUTO CDC FROM SNAPSHOT pour traiter les instantanés dans des tables cibles de type SCD 1 ou SCD 2. Pour plus d'informations sur l'utilisation de cette API, consultez Capture des données modifiées et instantanés.
Exemple : Traiter les instantanés en utilisant le temps d'ingestion du pipeline
Utilisez cette approche 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 le versioning. Un nouvel instantané est ingéré à chaque mise à jour du pipeline.
Vous pouvez lire les instantanés de plusieurs types de sources, y compris les tables Delta, les fichiers de stockage cloud et les connexions JDBC.
Étape 1 : Créer un échantillon de données
Créer une table contenant des données d'instantané. Exécutez le code suivant à partir d'un Notebook ou de Databricks SQL dans le dossier explorations de votre pipeline :
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.snapshot (
userId INT,
city STRING
);
INSERT INTO main.cdc_tutorial.snapshot VALUES
(1, 'Oaxaca'),
(2, 'Monterrey'),
(3, 'Tijuana');
Étape 2 : exécutez AUTO CDC FROM SNAPSHOT
Qu'est-ce qu'un pipeline ? pour exécuter le code dans cette étape.
Choisissez un type de source pour la vue d'instantané (le code de création d'exemple génère une table Delta) :
Option A : lire à partir d'une table Delta
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.table("main.cdc_tutorial.snapshot")
Option B : Lecture depuis le stockage cloud
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.format("csv").option("header", True).load("<snapshot-path>")
Option C : Lecture depuis JDBC (compute classique uniquement)
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.load()
)
Toutes les options, écrire dans la cible
Ensuite, ajoutez la table cible et le flux :
dp.create_streaming_table("target")
dp.create_auto_cdc_from_snapshot_flow(
target = "target",
source = "source",
keys = ["userId"],
stored_as_scd_type = 2
)
Après la 1re exécution du pipeline, tous les enregistrements sont insérés comme des lignes actives :
ID utilisateur | ville | start | __END_AT |
|---|---|---|---|
1 | Oaxaca | 0 | nul |
2 | Monterrey | 0 | nul |
3 | Tijuana | 0 | nul |
Pour utiliser le SCD de type 1 et ne conserver que l'état actuel, définissez stored_as_scd_type=1. Dans ce cas, la table cible n'inclut pas les colonnes __START_AT et __END_AT.
Étape 3 : Simulez un nouvel instantané et réexécutez.
Mettez à jour la table source pour simuler l'arrivée d'un nouvel instantané (exécutez ce code à partir d'un notebook ou d'un fichier SQL dans le dossier explorations de votre pipeline) :
TRUNCATE TABLE main.cdc_tutorial.snapshot;
INSERT INTO main.cdc_tutorial.snapshot VALUES
(2, 'Carmel'),
(3, 'Los Angeles'),
(4, 'Death Valley'),
(6, 'Kings Canyon');
Exécutez à nouveau le pipeline. AUTO CDC FROM SNAPSHOT compare le nouvel instantané au précédent et détecte que l'utilisateur 1 a été supprimé, les utilisateurs 2 et 3 ont été mis à jour et les utilisateurs 4 et 6 ont été insérés. Ceci génère un flux de modifications et utilise AUTO CDC pour créer la table de sortie.
Après la 2e exécution avec le SCD de type 2, la table cible contient les enregistrements suivants :
ID utilisateur | ville | start | __END_AT |
|---|---|---|---|
1 | Oaxaca | 0 | 1 |
2 | Monterrey | 0 | 1 |
2 | Carmel | 1 | nul |
3 | Tijuana | 0 | 1 |
3 | Los Angeles | 1 | nul |
4 | Death Valley | 1 | nul |
6 | Kings Canyon | 1 | nul |
L'utilisateur 1 a été supprimé. Les utilisateurs 2 et 3 ont chacun deux versions montrant leurs changements de ville. Les utilisateurs 4 et 6 ont été nouvellement insérés.
Après la deuxième exécution avec SCD de type 1, la table cible affiche uniquement l'état actuel :
ID utilisateur | ville |
|---|---|
2 | Carmel |
3 | Los Angeles |
4 | Death Valley |
6 | Kings Canyon |
Exemple : Traiter les instantanés à l'aide de fonctions de version
Utilisez cette approche lorsque vous avez besoin d'un contrôle explicite sur l'ordonnancement des instantanés. Par exemple, utilisez cette approche lorsque plusieurs instantanés arrivent en même temps, ou lorsque les instantanés arrivent dans le désordre. Vous écrivez une fonction qui spécifie quel instantané traiter ensuite ainsi que son numéro de version. L'API traite les instantanés par ordre de version croissant :
- Si plusieurs instantanés sont en stockage, ils sont tous traités dans l'ordre.
- Si un instantané arrive dans le désordre (par exemple, si
snapshot_3arrive aprèssnapshot_4), il est ignoré. - S'il n'y a pas de nouveaux snapshots, la fonction renvoie
Noneet aucun traitement n'a lieu.
Étape 1 : Préparez les fichiers de snapshot
Créez des fichiers CSV contenant des données d'instantané et ajoutez-les à un volume ou à un emplacement de stockage cloud. Nommez les fichiers de manière chronologique (par exemple, snapshot_1.csv, snapshot_2.csv).
Chaque fichier doit contenir des colonnes pour userId et city. Par exemple :
snapshot_1.csv :
ID utilisateur | ville |
|---|---|
1 | Oaxaca |
2 | Monterrey |
3 | Tijuana |
snapshot_2.csv :
ID utilisateur | ville |
|---|---|
2 | Carmel |
3 | Los Angeles |
4 | Death Valley |
Étape 2 : Exécutez AUTO CDC FROM SNAPSHOT avec une fonction de version
Créez un nouveau Notebook et collez le code pipeline suivant. Ensuite, Qu'est-ce qu'un pipeline ?.
from pyspark import pipelines as dp
from typing import Optional, Tuple
from pyspark.sql import DataFrame
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Optional[Tuple[DataFrame, int]]:
snapshot_dir = "/Volumes/main/cdc_tutorial/snapshots/" # or the location you created the sample data
files = dbutils.fs.ls(snapshot_dir)
snapshot_files = [f.name for f in files if f.name.startswith("snapshot_") and f.name.endswith(".csv")]
snapshot_versions = []
for filename in snapshot_files:
try:
version = int(filename.replace("snapshot_", "").replace(".csv", ""))
snapshot_versions.append(version)
except ValueError:
continue
snapshot_versions.sort()
if latest_snapshot_version is None:
if snapshot_versions:
next_version = snapshot_versions[0]
else:
return None
else:
next_versions = [v for v in snapshot_versions if v > latest_snapshot_version]
if next_versions:
next_version = next_versions[0]
else:
return None
snapshot_path = f"{snapshot_dir}snapshot_{next_version}.csv"
df = spark.read.format("csv").option("header", True).load(snapshot_path)
return (df, next_version)
dp.create_streaming_table("main.cdc_tutorial.target_versioned")
dp.create_auto_cdc_from_snapshot_flow(
target = "main.cdc_tutorial.target_versioned",
source = next_snapshot_and_version,
keys = ["userId"],
stored_as_scd_type = 2
)
Pour utiliser le type SCD 1 plutôt, définissez stored_as_scd_type=1.
Après traitement de snapshot_1.csv, la table cible contient les enregistrements suivants :
ID utilisateur | ville | start | __END_AT |
|---|---|---|---|
1 | Oaxaca | 1 | nul |
2 | Monterrey | 1 | nul |
3 | Tijuana | 1 | nul |
Après traitement de snapshot_2.csv, la table cible contient les enregistrements suivants :
ID utilisateur | ville | start | __END_AT |
|---|---|---|---|
1 | Oaxaca | 1 | 2 |
2 | Monterrey | 1 | 2 |
2 | Carmel | 2 | nul |
3 | Tijuana | 1 | 2 |
3 | Los Angeles | 2 | nul |
4 | Death Valley | 2 | nul |
N'oubliez pas que, pour le SCD de type 1, la table ressemble exactement à l'instantané le plus récent. La différence est que les queries en aval peuvent utiliser le flux de modifications pour ne traiter que les enregistrements modifiés.
Étape 3 : ajouter de nouveaux instantanés
Ajoutez un nouveau fichier CSV à l'emplacement de stockage avec des données modifiées (par exemple, des valeurs de ville modifiées, de nouvelles lignes ou des lignes supprimées). Exécutez ensuite le pipeline à nouveau pour traiter le nouvel instantané.
Limitations
- La colonne de séquencement doit être un type de données triable. Les valeurs de séquencement
NULLne sont pas prises en charge. AUTO CDC FROM SNAPSHOTn'est pris en charge que dans l'interface de pipeline Python ; l'interface SQL n'est pas prise en charge.- Pour Stream des données depuis la cible d'un processus AUTO CDC, lisez à partir de son flux de changements. Pour plus de détails, consultez Lire un flux de données de modification d'une table cible AUTO CDC.
Ressources supplémentaires
- Capture des changements de données et instantanés: Apprenez-en davantage sur les concepts de la CDC, les instantanés et les types de SCD.
- Répliquer une table RDBMS externe à l'aide de
AUTO CDC: Découvrez comment effectuer une hydratation initiale avec un pipelineonce, puis continuer à traiter les modifications. - Rubriques avancées sur AUTO CDC: découvrez les Opérations de changement sur les cibles AUTO CDC, la lecture des flux de données de changement et le traitement des métriques.
- Tutoriel : Construire un pipeline ETL en utilisant la capture de données modifiées