create_auto_cdc_from_snapshot_flow
La fonction create_auto_cdc_from_snapshot_flow crée un flux qui utilise la fonctionnalité de capture des données modifiées (CDC) de Lakeflow Pipelines pour traiter les données sources à partir d’instantanés de base de données. Consultez Fonctionnement d’AUTO CDC FROM SNAPSHOT.
Cette fonction remplace la fonction précédente apply_changes_from_snapshot(). Les deux fonctions ont la même signature. Databricks vous recommande d’effectuer la mise à jour pour utiliser le nouveau nom.
Vous devez disposer d'une table de streaming cible pour cette opération. Pour créer la table cible requise, vous pouvez utiliser la fonction create_streaming_table(). Vous ne pouvez pas cibler la même table de streaming avec create_auto_cdc_from_snapshot_flow() et create_auto_cdc_flow().
Syntaxe
from pyspark import pipelines as dp
dp.create_auto_cdc_from_snapshot_flow(
target = "<target-table>",
source = Any,
keys = ["key1", "key2", "keyN"],
stored_as_scd_type = "1",
track_history_column_list = None,
track_history_except_column_list = None
)
Pour le traitement AUTO CDC FROM SNAPSHOT, le comportement default est d'insérer une nouvelle ligne lorsqu'un enregistrement correspondant avec la même(les mêmes) clé(s) n'existe pas dans la cible. Si un enregistrement correspondant existe, il est mis à jour uniquement si l'une des valeurs de la ligne a changé. Les lignes avec des clés présentes dans la cible mais qui ne sont plus présentes dans la source sont supprimées.
Pour en savoir plus sur le traitement CDC avec des instantanés, consultez Les API AUTO CDC : simplifier la capture des modifications de données avec des pipelines. Pour des exemples d'utilisation de la fonction create_auto_cdc_from_snapshot_flow(), consultez les exemples d'ingestion instantanée périodique et d'ingestion instantanée historique.
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. Soit le nom d'une table ou d'une vue à instantanéer périodiquement, soit une fonction lambda Python qui renvoie le DataFrame d'instantané à traiter et la version de l'instantané. Voir la mise en œuvre de l'argument |
|
| 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 :
Les arguments des fonctions |
|
| Indique si les enregistrements doivent être stockés en tant que SCD de type 1 ou SCD de type 2. Définissez sur |
|
| Un sous-ensemble de colonnes de sortie à suivre pour l’historique dans la table cible. Utilisez
Les arguments des fonctions |
Implémentez l'argument source
La fonction create_auto_cdc_from_snapshot_flow() inclut l'argument source. Pour le traitement des instantanés historiques, l'argument source doit être une fonction lambda Python qui renvoie deux valeurs à la fonction create_auto_cdc_from_snapshot_flow() : un DataFrame Python contenant les données d'instantané à traiter et une version d'instantané.
Voici la signature de la fonction lambda :
lambda Any => Optional[(DataFrame, Any)]
- L'argument de la fonction lambda est la version d'instantané la plus récemment traitée.
- La valeur de retour de la fonction lambda est
Noneou un tuple de deux valeurs : la première valeur du tuple est un DataFrame contenant l'instantané à traiter. La deuxième valeur du tuple est la version de l'instantané qui représente l'ordre logique de l'instantané.
Un exemple qui implémente et appelle la fonction lambda :
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Tuple[DataFrame, Optional[int]]:
if latest_snapshot_version is None:
return (spark.read.load("filename.csv"), 1)
else:
return None
create_auto_cdc_from_snapshot_flow(
# ...
source = next_snapshot_and_version,
# ...
)
L'environnement d'exécution des LakeFlow Pipelines effectue les étapes suivantes chaque fois que le pipeline qui contient la fonction create_auto_cdc_from_snapshot_flow() est déclenché :
- Exécute la fonction
next_snapshot_and_versionpour charger le DataFrame d’instantané suivant et la version d’instantané correspondante. - Si aucun DataFrame n'est renvoyé, l'exécution est terminée et la mise à jour du pipeline est marquée comme terminée.
- Détecte les changements dans le nouvel instantané et les applique de manière incrémentielle à la table cible.
- Retourne à l'étape n°1 pour charger le prochain instantané et sa version.