create_auto_cdc_flow
La fonction create_auto_cdc_flow() crée un flux qui utilise la fonctionnalité de capture de données modifiées (CDC) des LakeFlow Pipelines pour traiter les données sources d'un flux de données modifiées (CDF).
Cette fonction remplace la fonction précédente apply_changes(). Les deux fonctions ont la même signature. Databricks vous recommande d’effectuer la mise à jour pour utiliser le nouveau nom.
Vous devez déclarer une table de streaming cible pour y appliquer les modifications. Vous pouvez éventuellement spécifier le schéma de votre table cible. Lorsque vous spécifiez le schéma de la table cible create_auto_cdc_flow(), vous devez inclure les colonnes __START_AT et __END_AT avec le même type de données que les champs sequence_by.
Pour créer la table cible requise, vous pouvez utiliser la fonction create_streaming_table() dans l'interface Python du pipeline.
Syntaxe
from pyspark import pipelines as dp
dp.create_auto_cdc_flow(
target = "<target-table>",
source = "<data-source>",
keys = ["key1", "key2", "keyN"],
sequence_by = "<sequence-column>",
system_sequence_by = None, # optional
ignore_null_updates = False, # optional
ignore_null_updates_column_list = None, # optional
ignore_null_updates_except_column_list = None, # optional
columns_to_update = None, # optional
apply_as_deletes = None, # optional
apply_as_truncates = None, # optional
column_list = None, # optional
except_column_list = None, # optional
stored_as_scd_type = "1", # optional
track_history_column_list = None, # optional
track_history_except_column_list = None, # optional
name = None, # optional
once = False # optional
)
Pour create_auto_cdc_flow le traitement, le comportement par default pour les INSERT UPDATE événements et consiste à *upserter* les événements CDC à partir de la source : mettre à jour toutes les lignes de la table cible qui correspondent à la ou aux clés spécifiées ou insérer une nouvelle ligne lorsqu'un enregistrement correspondant n'existe pas dans la table cible. La gestion des événements DELETE peut être spécifiée avec le paramètre apply_as_deletes.
Pour en savoir plus sur le traitement CDC avec un flux de changements, consultez les APIs AUTO CDC : Simplifier la capture des données modifiées avec des pipelines. Pour un exemple d'utilisation de la fonction create_auto_cdc_flow(), consultez les exemples AUTO CDC.
parameter
parameter | Type | Description |
|---|---|---|
|
| Obligatoire. Le nom de la table à mettre à jour. Vous pouvez utiliser la fonction create_streaming_table() pour créer la table cible avant d'exécuter la fonction |
|
| Obligatoire. La source de données contenant des enregistrements CDC. |
|
| Obligatoire. La colonne ou la combinaison de colonnes qui identifie de manière unique une ligne dans les données source. Ceci est utilisé pour identifier quels événements CDC s'appliquent à des enregistrements spécifiques dans la table cible. Vous pouvez spécifier au choix :
|
|
| Obligatoire. Les noms de colonnes spécifiant l'ordre logique des événements CDC dans les données source. LakeFlow Pipelines utilisent ce séquençage pour gérer les événements de modification qui arrivent dans le désordre. La colonne spécifiée doit être un type de données triable. Vous pouvez spécifier :
|
|
| La colonne spécifiant l'heure système à laquelle chaque événement CDC est connu du système. Utilisé avec Ce paramètre est facultatif et s’applique uniquement aux tables bitemporelles. |
|
| Contrôle la façon dont les valeurs Défini sur La default est Pour un contrôle plus précis des colonnes qui ignorent les valeurs |
|
| Un sous-ensemble de colonnes pour lesquelles les valeurs |
|
| Un sous-ensemble de colonnes qui appliquent des valeurs |
|
| Le nom d'une colonne source qui contient, pour chaque enregistrement de modification, l'ensemble des colonnes à mettre à jour sous forme de tableau de chaînes de noms de colonne ( Vous ne pouvez pas définir |
|
| Spécifie quand un événement CDC doit être traité comme un
Pour gérer les données désordonnées, la ligne supprimée est temporairement conservée comme marque de suppression dans la table Delta sous-jacente, et une vue est créée dans le métastore qui filtre ces marques de suppression. L'intervalle de rétention est de deux jours par default et peut être configuré avec la propriété de table Si vous utilisez Auto Loader comme source pour votre pipeline CDC, Auto Loader ne garantit pas l'ordre de traitement des fichiers. Pour plus de détails, consultez Gérer les données désordonnées. Définissez |
|
| Spécifie quand un événement CDC doit être traité comme une table complète
Étant donné que cette clause Trigger une troncation complète de la table cible, elle ne doit être utilisée que pour des cas d'utilisation spécifiques nécessitant cette fonctionnalité. Le |
|
| Un sous-ensemble de colonnes à inclure dans la table cible. Utilisez
Les arguments des fonctions |
|
| Pour stocker des enregistrements en tant que type SCD 1, type SCD 2 ou bitemporels. Définissez sur |
|
| Un sous-ensemble de colonnes de sortie à suivre pour l’historique dans la table cible. Utilisez
Les arguments des fonctions |
|
| Le nom du flux. Si non spécifié, default est la même valeur que |
|
| Vous pouvez éventuellement définir le flux comme un flux ponctuel, tel qu’un remplissage rétroactif. L'utilisation de
|
Notes
Pour les cibles SCD de type 1 dans les pipelines déclenchés, un ou plusieurs flux create_auto_cdc_flow() peuvent partager une cible avec un flux AUTO CDC FROM SNAPSHOT défini en Python ou en SQL. Les conditions suivantes s'appliquent :
- Attribuez à chaque
create_auto_cdc_flow()unnameunique. - Utilisez le même nombre de clés dans le même ordre pour tous les flux. Les noms de clés du flux d'instantanés sont comparés sans distinction de majuscules et de minuscules avec les noms de clés
AUTO CDC. Plusieurs fluxcreate_auto_cdc_flow()doivent utiliser des noms de clés et des casses identiques. - Utilisez exactement le même type de données pour
sequence_byet la version d’instantané. - Ne définissez pas d’attentes sur le flux
AUTO CDC FROM SNAPSHOT. - Définissez
ignore_null_updatessurFalse, et ne définissez niignore_null_updates_column_list, niignore_null_updates_except_column_list. - Si vous ajoutez un flux
AUTO CDCà une cibleAUTO CDC FROM SNAPSHOTexistante, la cible ne doit pas contenir de colonnes utilisateur dont les noms sont en conflit avec les colonnes systèmeAUTO CDCréservées. - N'utilisez pas ce modèle pour les cibles SCD Type 2 ou bitemporelles.
For an example, see Add a backfill to an AUTO CDC SCD Type 1 table.