Aller au contenu principal

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.

remarque

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.

important

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

Python
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
)
remarque

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

target

str

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 create_auto_cdc_from_snapshot_flow().

source

str OU lambda function

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 source.

keys

list

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 :

  • Une liste de chaînes de caractères : ["userId", "orderId"]

  • Une liste de fonctions Spark SQL col() : [col("userId"), col("orderId"].

Les arguments des fonctions col() ne peuvent pas inclure de qualificateurs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId).

stored_as_scd_type

str OU int

Indique si les enregistrements doivent être stockés en tant que SCD de type 1 ou SCD de type 2. Définissez sur 1 pour le SCD de type 1 ou 2 pour le SCD de type 2. Le default est SCD de type 1.

track_history_column_list OU track_history_except_column_list

list

Un sous-ensemble de colonnes de sortie à suivre pour l’historique dans la table cible. Utilisez track_history_column_list pour spécifier la liste complète des colonnes à suivre. Utilisez track_history_except_column_list pour spécifier les colonnes à exclure du suivi. Vous pouvez déclarer chaque valeur comme une liste de chaînes de caractères ou comme des fonctions col() Spark SQL :

  • track_history_column_list = ["userId", "name", "city"]
  • track_history_column_list = [col("userId"), col("name"), col("city")]
  • track_history_except_column_list = ["operation", "sequenceNum"]
  • track_history_except_column_list = [col("operation"), col("sequenceNum")

Les arguments des fonctions col() ne peuvent pas inclure de qualificateurs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId). The default is to inclure toutes les colonnes dans la table cible lorsqu'aucun argument track_history_column_list ou track_history_except_column_list n'est transmis à la fonction.

parameter

Type

Description

target

str

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 create_auto_cdc_from_snapshot_flow().

source

str OU lambda function

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 source.

keys

list

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 :

  • Une liste de chaînes de caractères : ["userId", "orderId"]

  • Une liste de fonctions Spark SQL col() : [col("userId"), col("orderId"].

Les arguments des fonctions col() ne peuvent pas inclure de qualificateurs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId).

stored_as_scd_type

str OU int

Indique si les enregistrements doivent être stockés en tant que SCD de type 1 ou SCD de type 2. Définissez sur 1 pour le SCD de type 1 ou 2 pour le SCD de type 2. Le default est SCD de type 1.

track_history_column_list OU track_history_except_column_list

list

Un sous-ensemble de colonnes de sortie à suivre pour l’historique dans la table cible. Utilisez track_history_column_list pour spécifier la liste complète des colonnes à suivre. Utilisez track_history_except_column_list pour spécifier les colonnes à exclure du suivi. Vous pouvez déclarer chaque valeur comme une liste de chaînes de caractères ou comme des fonctions col() Spark SQL :

  • track_history_column_list = ["userId", "name", "city"]
  • track_history_column_list = [col("userId"), col("name"), col("city")]
  • track_history_except_column_list = ["operation", "sequenceNum"]
  • track_history_except_column_list = [col("operation"), col("sequenceNum")

Les arguments des fonctions col() ne peuvent pas inclure de qualificateurs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId). The default is to inclure toutes les colonnes dans la table cible lorsqu'aucun argument track_history_column_list ou track_history_except_column_list n'est transmis à la fonction.

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 :

Python
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 None ou 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 :

Python
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é :

  1. Exécute la fonction next_snapshot_and_version pour charger le DataFrame d’instantané suivant et la version d’instantané correspondante.
  2. Si aucun DataFrame n'est renvoyé, l'exécution est terminée et la mise à jour du pipeline est marquée comme terminée.
  3. Détecte les changements dans le nouvel instantané et les applique de manière incrémentielle à la table cible.
  4. Retourne à l'étape n°1 pour charger le prochain instantané et sa version.