Rubriques avancées AUTO CDC
Au-delà des API AUTO CDC et AUTO CDC FROM SNAPSHOT de base, vous pouvez exécuter des DML sur les tables cibles, lire les flux de données de modification à partir des cibles CDC, surveiller les métriques de traitement, appliquer des mises à jour partielles et suivre les modifications avec un stockage bitemporel. Pour une introduction aux AUTO CDC APIs, consultez Les APIs AUTO CDC : simplifier la capture des modifications de données avec des pipelines.
Ajouter, modifier ou supprimer des données dans une table de streaming cible
Si votre pipeline publie des tables dans Unity Catalog, vous pouvez utiliser des instructions de langage de manipulation de données (LMD), y compris les instructions d'insertion, de mise à jour, de suppression et de Merge, pour modifier les tables de streaming cibles créées par les instructions AUTO CDC ... INTO.
- Les instructions DML qui modifient le schéma de table d'une table en streaming ne sont pas prises en charge. Assurez-vous que vos instructions DML ne tentent pas de faire évoluer le schéma de la table.
- Les instructions DML qui mettent à jour une table de streaming peuvent être exécutées uniquement dans un cluster Unity Catalog partagé ou un SQL Warehouse utilisant Databricks Runtime 13.3 LTS et versions ultérieures.
- Étant donné que le streaming nécessite des sources de données d'ajout uniquement, si votre traitement nécessite de lire en streaming une table de streaming source avec des modifications (par exemple, par des instructions DML), définissez l'indicateur skipChangeCommits lors de la lecture de la table de streaming source. Lorsque
skipChangeCommitsest défini, les transactions qui suppriment ou modifient des enregistrements sur la table source sont ignorées. Si votre traitement ne nécessite pas de table de streaming, vous pouvez utiliser une vue matérialisée (qui n'a pas la restriction d'ajout uniquement) comme table cible.
Comme le pipeline utilise une colonne SEQUENCE BY spécifiée et propage les valeurs de séquencement appropriées aux colonnes __START_AT et __END_AT de la table cible (pour SCD de type 2), vous devez vous assurer que les instructions DML utilisent des valeurs valides pour ces colonnes afin de maintenir l'ordre correct des enregistrements. Découvrir le fonctionnement d’AUTO CDC.
Pour plus d'informations sur l'utilisation des instructions DML avec les tables de streaming, consultez Ajouter, modifier ou supprimer des données dans une table de streaming.
L'exemple suivant insère un enregistrement actif avec une start sequence de 5 :
INSERT INTO my_streaming_table (id, name, __START_AT, __END_AT) VALUES (123, 'John Doe', 5, NULL);
Si vous devez renommer les colonnes __START_AT et __END_AT dans votre table cible SCD de type 2 (par exemple, pour correspondre aux exigences de schéma en aval), créez une vue sur la table cible :
CREATE VIEW my_employees_view AS
SELECT
*,
__START_AT AS valid_from,
__END_AT AS valid_to
FROM my_scd2_target_table;
Lire un flux de données modifiées à partir d'une table cible AUTO CDC
Dans Databricks Runtime 15.2 et versions ultérieures, vous pouvez lire un flux de données de modification à partir d'une table de streaming qui est la cible de requêtes AUTO CDC ou AUTO CDC FROM SNAPSHOT de la même manière que vous lisez un flux de données de modification à partir d'autres tables Delta. Les éléments suivants sont requis pour lire le flux de données de modification à partir d'une table de streaming cible :
- La table de streaming cible doit être publiée dans Unity Catalog. Consultez Utiliser Unity Catalog avec des pipelines.
- Pour lire le flux de données de modification de la table de streaming cible, vous devez utiliser Databricks Runtime 15.2 ou version ultérieure. Pour lire le flux de données de modification dans un autre pipeline, le pipeline doit être configuré pour utiliser Databricks Runtime 15.2 ou version ultérieure.
Vous lisez le flux de données de modification à partir d'une table de streaming cible créée dans un Lakeflow Pipelines de la même manière que la lecture d'un flux de données de modification à partir d'autres tables Delta. Pour en savoir plus sur l'utilisation de la fonctionnalité de flux de données de modification Delta, y compris des exemples en Python et SQL, consultez Utiliser le flux de données de modification sur Databricks.
L'enregistrement du flux de données modifiées inclut des métadonnées identifiant le type d'événement de modification. Lorsqu'un enregistrement est mis à jour dans une table, les métadonnées des enregistrements de modification associés incluent généralement les valeurs _change_type définies sur update_preimage et les événements update_postimage.
Cependant, les valeurs _change_type sont différentes si des mises à jour sont effectuées sur la table de streaming cible, incluant la modification des valeurs de clé primaire. Lorsque les modifications incluent des mises à jour des clés primaires, les champs de métadonnées _change_type sont définis sur les événements insert et delete. Des modifications des clés primaires peuvent se produire lorsque des mises à jour manuelles sont effectuées sur l'un des champs clés avec une instruction UPDATE ou MERGE ou, pour les tables de type SCD 2, lorsque le champ __start_at change pour refléter une valeur de séquence de start antérieure.
La query AUTO CDC détermine les valeurs de la clé primaire, qui diffèrent pour le traitement SCD de type 1 et SCD de type 2 :
SCD Type | Primary key |
|---|---|
SCD type 1, and the pipelines Python interface | The primary key is the value of the |
SCD type 2 | The primary key is the |
Obtenir des données sur les enregistrements traités par une query CDC dans les pipelines
Les métriques suivantes sont capturées uniquement par AUTO CDC queries et non par AUTO CDC FROM SNAPSHOT queries.
Les métriques suivantes sont capturées par AUTO CDC queries :
num_upserted_rows: Le nombre de lignes de sortie insérées dans le dataset lors d'une mise à jour.num_deleted_rows: Le nombre de lignes de sortie existantes supprimées du dataset lors d'une mise à jour.
La métrique num_output_rows, sortie pour les flux non-CDC, n'est pas capturée pour AUTO CDC requêtes.
Appliquer des mises à jour partielles
Lorsqu'une source n'envoie que les colonnes qui ont changé, AUTO CDC doit faire la distinction entre une colonne absente d'un enregistrement de modification, qui devrait laisser la valeur cible inchangée, et une colonne explicitement définie sur null, qui devrait écraser la valeur cible avec null. By default, IGNORE NULL UPDATES traite chaque null comme un marqueur « ne pas mettre à jour », il ne peut donc pas appliquer un null explicite. Pour résoudre cette ambiguïté, choisissez l'une des trois méthodes suivantes :
Méthode | Quand utiliser | Comportement |
|---|---|---|
| Un petit ensemble fixe de colonnes devrait ignorer les valeurs | Les colonnes répertoriées conservent leur valeur cible existante lorsque la valeur entrante est |
| La plupart des colonnes doivent ignorer les valeurs | Les colonnes spécifiées appliquent des valeurs |
| Chaque enregistrement de modification met à jour un ensemble différent de colonnes, ou l'ensemble des colonnes modifiables change au fil du temps. | Une colonne source nomme les colonnes à mettre à jour pour chaque enregistrement de modification. Les colonnes listées sont écrites à partir de la source, y compris les valeurs |
COLUMNS TO UPDATE ne peut pas être combiné avec IGNORE NULL UPDATES et n'est pas pris en charge pour les tables bitemporelles.
En règle générale, choisissez COLUMNS TO UPDATE lorsque le producteur sait quelles colonnes ont été modifiées dans chaque enregistrement et peut transporter cette information dans une colonne source, par exemple lorsque plusieurs producteurs écrivent à la même source ou que l'ensemble des colonnes modifiables augmente avec le temps. Choisissez IGNORE NULL UPDATES ON lorsque le propriétaire du pipeline connaît à l’avance l’ensemble fixe des colonnes actualisables et préfère les contrôler dans le code du pipeline.
L'exemple suivant utilise une colonne source nommée columnsToUpdate pour contrôler les colonnes que chaque enregistrement de modification met à jour, y compris les colonnes définies explicitement sur null:
- Python
- SQL
from pyspark import pipelines as dp
dp.create_streaming_table("target")
dp.create_auto_cdc_flow(
target = "target",
source = "cdc_source",
keys = ["id"],
sequence_by = "sequenceNum",
stored_as_scd_type = 1,
columns_to_update = "columnsToUpdate"
)
CREATE OR REFRESH STREAMING TABLE target;
CREATE FLOW apply_cdc AS AUTO CDC INTO
target
FROM
stream(cdc_source)
KEYS
(id)
SEQUENCE BY
sequenceNum
STORED AS
SCD TYPE 1
COLUMNS TO UPDATE
columnsToUpdate;
Pour la référence complète des parameters, consultez AUTO CDC INTO (pipelines) et create_auto_cdc_flow.
AUTO CDC bitemporel
Bêta
Le CDC AUTO bitemporel est en bêta.
Les SCD de type 1 et de type 2 sont unitemporels : ils suivent les changements sur une seule dimension temporelle. Le bitemporel étend l'historique du SCD de type 2 pour suivre les modifications sur deux dimensions temporelles et distinguer deux perspectives :
- Heure d'activité : lorsque l'événement s'est réellement produit.
- Heure système : lorsque le système a enregistré ou ingéré l’événement.
À l'instar du SCD de Type 2, la bitemporalité préserve un historique complet des enregistrements. Il ajoute une deuxième chronologie afin que vous puissiez reconstituer à la fois ce que les données montraient et ce que le système croyait à tout moment par le passé.
Par exemple, un fonds spéculatif ingère des données boursières depuis un système source. Le cours de l'action d'Acme Corp change le 1er janvier, mais le fonds n'ingère cette mise à jour que le 5 janvier. La CDC AUTO bitemporelle permet au fonds de répondre à deux questions distinctes : quel était le cours réel de l'action d'Acme Corp le 1er janvier (heure d'activité), et quel était le prix que le système croyait lorsque le fonds a pris des décisions de trading le 3 janvier (heure système). La capacité à distinguer ces chronologies est utile pour l'audit, les rapports réglementaires et la prise de décision financière.
Pour activer le traitement bitemporel, définissez STORED AS BITEMPORAL (SQL) ou stored_as_scd_type="bitemporal" (Python), utilisez SEQUENCE BY pour la colonne de temps métier et SYSTEM SEQUENCE BY pour la colonne de temps système. La table cible ajoute les colonnes __SYSTEM_START_AT et __SYSTEM_END_AT en plus des colonnes SCD de type 2 __START_AT et __END_AT. Pour les détails de la syntaxe, consultez AUTO CDC INTO (pipelines) ou create_auto_cdc_flow.
Exemples de CDC AUTO bitemporels
L’exemple suivant crée une table cible bitemporelle à partir d’un petit ensemble d’événements CDC synthétiques. La colonne bt contient l’heure métier et la colonne st contient l’heure système.
- Python
- SQL
from pyspark import pipelines as dp
# Source: synthetic CDC events
dp.create_streaming_table(name="cdc_source")
@dp.append_flow(target="cdc_source", once=True)
def load_cdc_source():
return spark.createDataFrame(
[
(1, "x10", "y10", 10, 100),
(1, "x20", "y20", 20, 200)
],
schema="id INT, x STRING, y STRING, bt INT, st INT",
)
# Target: bitemporal table
dp.create_streaming_table(name="target_bitemporal")
dp.create_auto_cdc_flow(
target = "target_bitemporal",
source = "cdc_source",
keys = ["id"],
sequence_by = "bt",
system_sequence_by = "st",
stored_as_scd_type = "bitemporal"
)
-- Source: synthetic CDC events
CREATE OR REFRESH STREAMING TABLE cdc_source_sql;
CREATE FLOW cdc_source_sql AS INSERT INTO ONCE
cdc_source_sql BY NAME
SELECT * FROM VALUES
(1, 'x10', 'y10', 10, 100),
(1, 'x20', 'y20', 20, 200)
AS t(id, x, y, bt, st);
-- Target: bitemporal table
CREATE OR REFRESH STREAMING TABLE target_bitemporal_sql;
CREATE FLOW target_bitemporal_sql AS AUTO CDC INTO
target_bitemporal_sql
FROM
stream(cdc_source_sql)
KEYS
(id)
SEQUENCE BY
bt
SYSTEM SEQUENCE BY
st
STORED AS
BITEMPORAL;
La séquence de modifications suivante montre comment une table bitemporelle enregistre une insertion, une mise à jour, une mise à jour désordonnée et une suppression pour une seule entreprise. La colonne de séquençage génère les colonnes __START_AT et __END_AT (heure commerciale), et la colonne de séquençage système génère les colonnes __SYSTEM_START_AT et __SYSTEM_END_AT (heure système) :
Colonne | Description |
|---|---|
| L’heure d’activité à laquelle cette ligne est devenue valide. |
| L'heure commerciale à laquelle la validité de cette ligne se termine. |
| L'heure système à laquelle les données de cette ligne et l'intervalle de temps commercial sont réputés être vrais. |
| L'heure système à laquelle les données et l'intervalle de temps commercial de cette ligne sont connus pour être invalidés. |
Le système gère les événements qui arrivent dans n'importe quel ordre sur les deux chronologies. Lorsqu'un événement arrive avec une heure commerciale ou une heure système antérieure à celle des événements déjà traités, le système corrige l'historique affecté plutôt que de simplement l'ajouter à la fin.
Changement 1 : Insertion
L'entreprise A est ajoutée le 18/07/2025 à 10 h 01 min 00 s (heure ouvrée) mais n'est ingérée qu'à 10 h 05 min 00 s (heure système).
Entrée :
CompanyId | points de données | Séquençage | Séquençage système | Opérations |
|---|---|---|---|---|
A | XFv1 | 18/07/2025 10:01:00 | 18/07/2025 10:05:00 |
|
Résultat :
CompanyId | points de données | start | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
A | XFv1 | 18/07/2025 10:01:00 | NULL | 18/07/2025 10:05:00 | NULL |
XFv1 est valide à partir de 10:01:00 sans fin connue. Le système a pris connaissance de ce fait à l'heure système 10:05:00, sans fin connue.
Changement 2 : Mise à jour
La société A est mise à jour le 18/07/2025 à 12 h 15 min 43 s (heure ouvrable), et le système consomme l'événement à 12 h 20 min 00 s (heure système). Le système préserve à la fois ce qu'il pensait avant que la mise à jour ne soit connue et l'historique commercial corrigé après l'ingestion de la mise à jour.
Entrée :
CompanyId | points de données | Séquençage | Séquençage système | Opérations |
|---|---|---|---|---|
A | XFv2 | 18/07/2025 12 h 15 min 43 s | 18.07.2025 12:20:00 |
|
Résultat :
CompanyId | points de données | start | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
A | XFv1 | 18/07/2025 10:01:00 | NULL | 18/07/2025 10:05:00 | 18.07.2025 12:20:00 |
A | XFv1 | 18/07/2025 10:01:00 | 18/07/2025 12 h 15 min 43 s | 18.07.2025 12:20:00 | NULL |
A | XFv2 | 18/07/2025 12 h 15 min 43 s | NULL | 18.07.2025 12:20:00 | NULL |
XFv1 était considéré comme valide à partir de 10:01:00 sans fin connue, et le système a maintenu cette croyance de 10:05:00 jusqu'à 12:20:00. Il est désormais connu que XFv1 n'est valide que jusqu'à 12:15:43, un historique corrigé effectif à partir de l'heure système 12:20:00 sans fin connue. XFv2 est valide à partir de 12:15:43 sans fin connue, et a été appris à l'heure système 12:20:00.
Modification 3 : Mise à jour dans le désordre
Une mise à jour hors séquence arrive, indiquant que la société A a été mise à jour le 18/07/2025 à 12:05:00 (heure commerciale), mais elle n'est pas ingérée avant 12:25:00 (heure système). Lorsqu'une mise à jour arrive plus tard dans l'heure système mais avec une heure commerciale antérieure, le système corrige l'heure commerciale historique et préserve à la fois ce qu'il croyait avant la mise à jour hors séquence et l'historique corrigé.
Entrée :
CompanyId | points de données | Séquençage | Séquençage système | Opérations |
|---|---|---|---|---|
A | XFv3 | 18/07/2025 12 h 05 min 00 s | 18/07/2025 12 h 25 min 00 s |
|
Résultat :
CompanyId | points de données | start | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
A | XFv1 | 18/07/2025 10:01:00 | NULL | 18/07/2025 10:05:00 | 18.07.2025 12:20:00 |
A | XFv1 | 18/07/2025 10:01:00 | 18/07/2025 12 h 15 min 43 s | 18.07.2025 12:20:00 | 18/07/2025 12 h 25 min 00 s |
A | XFv1 | 18/07/2025 10:01:00 | 18/07/2025 12 h 05 min 00 s | 18/07/2025 12 h 25 min 00 s | NULL |
A | XFv3 | 18/07/2025 12 h 05 min 00 s | 18/07/2025 12 h 15 min 43 s | 18/07/2025 12 h 25 min 00 s | NULL |
A | XFv2 | 18/07/2025 12 h 15 min 43 s | NULL | 18.07.2025 12:20:00 | NULL |
XFv1 était considéré valide de 10 h 01 min 00 s à 12 h 15 min 43 s, et cette validité est désormais valable en temps système jusqu'à 12 h 25 min 00 s. La nouvelle mise à jour corrige la validité commerciale de XFv1 pour se terminer à 12:05:00, un historique corrigé effectif à partir de l'heure système 12:25:00. XFv3 est maintenant considéré comme valide de 12:05:00 à 12:15:43, une validité établie dans le temps système à partir de 12:25:00 sans fin connue.
Changement 4 : Supprimer
La société A est supprimée le 18/07/2025 12 h 30 min 00 s, et le système consomme l'événement à 12 h 30 min 00 s. Puisqu'une opération de suppression représente la fin de l'existence commerciale de l'entité, le système ne crée aucune ligne de remplacement. XFv2 apparaît en deux lignes, préservant une piste d'audit complète de la date à laquelle l'entreprise a cessé d'exister et de la date à laquelle le système a pris connaissance de la suppression.
Entrée :
CompanyId | points de données | Séquençage | Séquençage système | Opérations |
|---|---|---|---|---|
A | XFv2 | 18/07/2025 12:30:00 | 18/07/2025 12:30:00 |
|
Résultat :
CompanyId | points de données | start | __END_AT | __SYSTEM_START_AT | __SYSTEM_END_AT |
|---|---|---|---|---|---|
A | XFv1 | 18/07/2025 10:01:00 | NULL | 18/07/2025 10:05:00 | 18.07.2025 12:20:00 |
A | XFv1 | 18/07/2025 10:01:00 | 18/07/2025 12 h 15 min 43 s | 18.07.2025 12:20:00 | 18/07/2025 12 h 25 min 00 s |
A | XFv1 | 18/07/2025 10:01:00 | 18/07/2025 12 h 05 min 00 s | 18/07/2025 12 h 25 min 00 s | NULL |
A | XFv3 | 18/07/2025 12 h 05 min 00 s | 18/07/2025 12 h 15 min 43 s | 18/07/2025 12 h 25 min 00 s | NULL |
A | XFv2 | 18/07/2025 12 h 15 min 43 s | NULL | 18.07.2025 12:20:00 | 18/07/2025 12:30:00 |
A | XFv2 | 18/07/2025 12 h 15 min 43 s | 18/07/2025 12:30:00 | 18/07/2025 12:30:00 | NULL |
XFv2 était valide de 12:15:43 sans fin connue, et le système a maintenu cette croyance de 12:20:00 à 12:30:00. Une fois la suppression ingérée, XFv2 est valide uniquement jusqu'à 12:30:00, un historique corrigé prenant effet à partir de l'heure système 12:30:00.
Quels objets de données sont utilisés pour le traitement CDC dans un pipeline ?
Lorsque vous déclarez la table cible dans le Hive metastore, deux structures de données sont créées :
- Une vue utilisant le nom attribué à la table cible.
- Une table de support interne utilisée par le pipeline pour gérer le traitement CDC. Cette table est nommée en préfixant
__apply_changes_storage_au nom de la table cible.
Par exemple, si vous déclarez une table cible nommée dp_cdc_target, vous voyez une vue nommée dp_cdc_target et une table nommée __apply_changes_storage_dp_cdc_target dans le metastore. Interrogez la vue pour accéder aux données traitées. Ne modifiez pas la table sous-jacente directement.
Ces structures de données s'appliquent uniquement au traitement AUTO CDC, pas au traitement AUTO CDC FROM SNAPSHOT. Ils s'appliquent également uniquement au Hive metastore, et non au Unity Catalog.