Aller au contenu principal

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.

remarque

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 CDC traite 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.

remarque

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

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

ID utilisateur

Nom

ville

Opérations

sequenceNum

nul

nul

nul

TRONQUER

3

remarque

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.

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

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

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

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

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

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
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"]
)

É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

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 :

SQL
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

Python
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

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

Python
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 :

Python
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

ID utilisateur

ville

start

__END_AT

1

Oaxaca

0

nul

2

Monterrey

0

nul

3

Tijuana

0

nul

remarque

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

SQL
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

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

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_3 arrive après snapshot_4), il est ignoré.
  • S'il n'y a pas de nouveaux snapshots, la fonction renvoie None et 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

ID utilisateur

ville

1

Oaxaca

2

Monterrey

3

Tijuana

snapshot_2.csv :

ID utilisateur

ville

2

Carmel

3

Los Angeles

4

Death Valley

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

Python
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
)
remarque

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

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

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

remarque

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 NULL ne sont pas prises en charge.
  • AUTO CDC FROM SNAPSHOT n'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