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
|