Aller au contenu principal

Répliquer une table RDBMS externe à l'aide de AUTO CDC

Vous pouvez répliquer une table depuis un système de gestion de base de données relationnelles externe (RDBMS) vers Databricks en utilisant l'API AUTO CDC dans les pipelines. Vous apprendrez :

  • Modèles courants pour la configuration des sources.
  • Comment effectuer une copie complète unique des données existantes à l'aide d'un flux once.
  • Comment ingérer en continu de nouvelles modifications à l'aide d'un flux change.

Ce modèle est idéal pour créer des tables de dimensions à évolution lente (SCD) ou pour maintenir une table cible synchronisée avec un système d’enregistrement externe.

Avant de commencer

Ce guide suppose que vous avez accès aux datasets suivants à partir de votre source :

  • Un instantané complet de la table source dans le stockage cloud. Ce dataset est utilisé pour le chargement initial.
  • Un flux de changements continu, peuplé dans le même emplacement de stockage cloud (par exemple, en utilisant Debezium, Kafka ou la CDC basée sur les logs). Ce flux est l'entrée pour le processus AUTO CDC en cours.

Configurer les vues source

Tout d'abord, définissez deux vues sources pour remplir la table cible rdbms_orders à partir d'un chemin de stockage dans le cloud orders_snapshot_path. Les deux sont conçues comme des vues de streaming sur des données brutes dans le stockage cloud. L'utilisation de vues offre une plus grande efficacité, car les données n'ont pas besoin d'être écrites avant d'être utilisées dans le processus AUTO CDC.

  • La première vue source est un instantané complet (full_orders_snapshot).
  • Le second est un flux de modification continu (rdbms_orders_change_feed).

Les exemples de ce guide utilisent le stockage cloud comme source, mais vous pouvez utiliser n'importe quelle source prise en charge par les tables de streaming.

full_orders_snapshot()

Cette étape crée un pipeline avec une vue qui lit l'instantané complet initial des données de commandes.

L'exemple Python suivant :

  • Utilise spark.readStream avec Auto Loader (format("cloudFiles"))
  • Lit les fichiers JSON à partir d'un répertoire défini par orders_snapshot_path
  • Définit includeExistingFiles sur true pour garantir que les données historiques déjà présentes dans le chemin sont traitées.
  • Définit inferColumnTypes sur true pour déduire le schéma automatiquement
  • Retourne toutes les colonnes avec .select("\*")
Python
@dp.view()
def full_orders_snapshot():
return (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.includeExistingFiles", "true")
.option("cloudFiles.inferColumnTypes", "true")
.load(orders_snapshot_path)
.select("*")
)

rdbms_orders_change_feed()

Cette étape crée une deuxième vue qui lit les données de modification incrémentielles (par exemple, à partir des Logs CDC ou des tables de modification). Il lit les données de orders_cdc_path et suppose que des fichiers JSON de style CDC sont régulièrement déposés dans ce chemin.

Python
@dp.view()
def rdbms_orders_change_feed():
return (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.includeExistingFiles", "true")
.option("cloudFiles.inferColumnTypes", "true")
.load(orders_cdc_path)

Hydratation initiale (flux unique)

Maintenant que les sources sont configurées, la logique AUTO CDC Merge les deux sources dans une table de streaming cible. Tout d'abord, utilisez un flux AUTO CDC unique avec ONCE=TRUE pour copier le contenu complet de la table RDBMS dans une table de streaming. Ceci prépare la table cible avec des données historiques sans les rejouer lors de futures mises à jour.

Python
from pyspark import pipelines as dp

# Step 1: Create the target streaming table

dp.create_streaming_table("rdbms_orders")

# Step 2: Once Flow — Load initial snapshot of full RDBMS table

dp.create_auto_cdc_flow(
flow_name = "initial_load_orders",
once = True, # one-time load
target = "rdbms_orders",
source = "full_orders_snapshot", # e.g., ingested from JDBC into bronze
keys = ["order_id"],
sequence_by = "timestamp",
stored_as_scd_type = "1"
)

Le flux once ne s'exécute qu'une seule fois. Les nouveaux fichiers ajoutés à full_orders_snapshot après la création du pipeline sont ignorés.

important

Effectuer un refresh complet sur la table streaming rdbms_orders réexécute le flux once. Si les données de l'instantané initial dans le cloud storage ont été supprimées, cela entraîne une perte de données.

Flux de modification continu (flux de changements)

Après le chargement initial de l'instantané, utilisez un autre flux AUTO CDC pour ingérer en continu les modifications du flux CDC du RDBMS. Cela maintient votre table rdbms_orders à jour avec les insertions, les mises à jour et les suppressions.

Python
from pyspark import pipelines as dp

# Step 3: Change Flow — Ingest ongoing CDC stream from source system

dp.create_auto_cdc_flow(
flow_name = "orders_incremental_cdc",
target = "rdbms_orders",
source = "rdbms_orders_change_feed", # e.g., ingested from Kafka or Debezium
keys = ["order_id"],
sequence_by = "timestamp",
stored_as_scd_type = "1"
)

Considérations

Idempotence du remplissage

Un flux once ne se réexécute que lorsque la table cible est entièrement actualisée.

Multiples flux

Vous pouvez utiliser plusieurs flux de modifications pour Merge des corrections, des données arrivées tardivement ou des flux alternatifs, mais tous doivent partager un schéma et des clés.

refresh complète

Un full refresh sur la table de streaming rdbms_orders réexécute le flux once. Cela peut entraîner une perte de données si l'emplacement de stockage cloud initial a supprimé les données du snapshot initial.

Ordre d'exécution du flux

L'ordre d'exécution du flux n'a pas d'importance. Le résultat final est le même.

Idempotence du remplissage

Un flux once ne se réexécute que lorsque la table cible est entièrement actualisée.

Multiples flux

Vous pouvez utiliser plusieurs flux de modifications pour Merge des corrections, des données arrivées tardivement ou des flux alternatifs, mais tous doivent partager un schéma et des clés.

refresh complète

Un full refresh sur la table de streaming rdbms_orders réexécute le flux once. Cela peut entraîner une perte de données si l'emplacement de stockage cloud initial a supprimé les données du snapshot initial.

Ordre d'exécution du flux

L'ordre d'exécution du flux n'a pas d'importance. Le résultat final est le même.

Ressources supplémentaires

  • Connecteur SQL Server entièrement managé dans Lakeflow Connect