Aller au contenu principal

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

remarque

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.

important

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

Python
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

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

source

str

Obligatoire. La source de données contenant des enregistrements CDC.

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 qualificatifs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId).

sequence_by

str, col() ou struct()

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 :

  • Une chaîne : "sequenceNum"
  • Une fonction col() de Spark SQL : col("sequenceNum"). Les arguments des fonctions col() ne peuvent pas inclure de qualificatifs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId).
  • Un struct() combinant plusieurs colonnes pour briser les liens : struct("timestamp_col", "id_col"), il trie d'abord par le premier champ de structure, puis par le deuxième champ en cas d'égalité, et ainsi de suite.

system_sequence_by

str OU col()

La colonne spécifiant l'heure système à laquelle chaque événement CDC est connu du système. Utilisé avec stored_as_scd_type="bitemporal" pour suivre les changements à travers le temps métier (sequence_by) et l'heure système. La colonne spécifiée doit être un type de données triable. L’AUTO CDC bitemporelle est en bêta. Voir AUTO CDC bitemporelle.

Ce paramètre est facultatif et s’applique uniquement aux tables bitemporelles.

ignore_null_updates

bool

Contrôle la façon dont les valeurs null des mises à jour CDC entrantes sont gérées. Lorsque ignore_null_updates est True, les colonnes null d'une mise à jour entrante sont ignorées ; la valeur existante dans la ligne cible est préservée. Ceci s'applique également aux colonnes imbriquées avec des valeurs null. Lorsque ignore_null_updates est False, les colonnes null d'une mise à jour entrante écrasent les valeurs existantes dans la cible. default to False.

Défini sur True lorsque les événements source n'incluent que les colonnes modifiées, de sorte que les colonnes inchangées ne soient pas écrasées par null.

La default est False.

Pour un contrôle plus précis des colonnes qui ignorent les valeurs null, utilisez ignore_null_updates_column_list ou ignore_null_updates_except_column_list.

ignore_null_updates_column_list

list

Un sous-ensemble de colonnes pour lesquelles les valeurs null dans un enregistrement de modification entrant sont ignorées, de sorte que chacune de ces colonnes conserve sa valeur existante dans la cible. Les colonnes en dehors de la liste appliquent des valeurs null explicites. Utilisez ce paramètre pour appliquer des mises à jour partielles lorsque votre source n'envoie que les colonnes modifiées. Équivalent à la clause SQL IGNORE NULL UPDATES ON columnList. Utilisez soit ignore_null_updates_column_list, soit ignore_null_updates_except_column_list, pas les deux.

ignore_null_updates_except_column_list

list

Un sous-ensemble de colonnes qui appliquent des valeurs null explicites. Toutes les autres colonnes ignorent les valeurs null dans un enregistrement de modification entrant et conservent leur valeur existante dans la cible. Équivalent à la clause SQL IGNORE NULL UPDATES ON * EXCEPT (...). Utilisez ignore_null_updates_column_list ou ignore_null_updates_except_column_list, pas les deux.

columns_to_update

str OU col()

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 (array<string>). Les colonnes ne figurant pas dans le tableau conservent leurs valeurs cibles existantes, tandis que les colonnes répertoriées sont écrites à partir de la source, y compris les valeurs null explicites. Utilisez ce paramètre lorsque chaque enregistrement de modification met à jour un ensemble de colonnes différent et que vous devez appliquer des valeurs null explicites. Équivalent à la clause SQL COLUMNS TO UPDATE.

Vous ne pouvez pas définir columns_to_update avec ignore_null_updates, ignore_null_updates_column_list ou ignore_null_updates_except_column_list. columns_to_update n'est pas pris en charge pour les tables bitemporelles.

apply_as_deletes

str OU expr()

Spécifie quand un événement CDC doit être traité comme un DELETE plutôt qu'un upsert. Vous pouvez spécifier l'une des options suivantes :

  • Une chaîne : "Operation = 'DELETE'"
  • Une fonction expr() de Spark SQL : expr("Operation = 'DELETE'")

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 pipelines.cdc.tombstoneGCThresholdInSeconds.

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 pipelines.cdc.tombstoneGCThresholdInSeconds sur une valeur qui dépasse le délai maximal attendu entre l'arrivée de l'événement et l'exécution du pipeline. Cela garantit que les marqueurs de suppression sont conservés suffisamment longtemps pour gérer correctement les événements de suppression tardifs ou désordonnés.

apply_as_truncates

str OU expr()

Spécifie quand un événement CDC doit être traité comme une table complète TRUNCATE. Vous pouvez spécifier l'une des options suivantes :

  • Une chaîne : "Operation = 'TRUNCATE'"
  • Une fonction expr() de Spark SQL : expr("Operation = 'TRUNCATE'")

É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 apply_as_truncates parameter est pris en charge uniquement pour le type SCD 1. Le type SCD 2 ne prend pas en charge les Opérations de troncation.

column_list OU except_column_list

list

Un sous-ensemble de colonnes à inclure dans la table cible. Utilisez column_list pour spécifier la liste complète des colonnes à inclure. Utilisez except_column_list pour spécifier les colonnes à exclure. Vous pouvez déclarer chaque valeur comme une liste de chaînes de caractères ou comme des fonctions col() Spark SQL :

  • column_list = ["userId", "name", "city"]
  • column_list = [col("userId"), col("name"), col("city")]
  • except_column_list = ["operation", "sequenceNum"]
  • 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 column_list ou except_column_list n'est transmis à la fonction.

stored_as_scd_type

str OU int

Pour stocker des enregistrements en tant que type SCD 1, type SCD 2 ou bitemporels. Définissez sur 1 pour le type SCD 1, 2 pour le type SCD 2, ou "bitemporal" pour suivre les modifications à la fois pour l'heure métier et l'heure système. Le bitemporel nécessite system_sequence_by et est en Bêta. Consultez Bitemporal AUTO CDC. Le default est le type de SCD 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.

name

str

Le nom du flux. Si non spécifié, default est la même valeur que target.

once

bool

Vous pouvez éventuellement définir le flux comme un flux ponctuel, tel qu’un remplissage rétroactif. L'utilisation de once=True modifie le flux de deux manières :

  • La valeur de retour. streaming-query. doit être un DataFrame batch dans ce cas, et non un DataFrame en streaming.
  • Le flux s'exécute une seule fois par default. Si le pipeline est mis à jour avec un refresh complet, alors le flux ONCE s'exécute de nouveau pour recréer les données.

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

source

str

Obligatoire. La source de données contenant des enregistrements CDC.

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 qualificatifs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId).

sequence_by

str, col() ou struct()

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 :

  • Une chaîne : "sequenceNum"
  • Une fonction col() de Spark SQL : col("sequenceNum"). Les arguments des fonctions col() ne peuvent pas inclure de qualificatifs. Par exemple, vous pouvez utiliser col(userId), mais vous ne pouvez pas utiliser col(source.userId).
  • Un struct() combinant plusieurs colonnes pour briser les liens : struct("timestamp_col", "id_col"), il trie d'abord par le premier champ de structure, puis par le deuxième champ en cas d'égalité, et ainsi de suite.

system_sequence_by

str OU col()

La colonne spécifiant l'heure système à laquelle chaque événement CDC est connu du système. Utilisé avec stored_as_scd_type="bitemporal" pour suivre les changements à travers le temps métier (sequence_by) et l'heure système. La colonne spécifiée doit être un type de données triable. L’AUTO CDC bitemporelle est en bêta. Voir AUTO CDC bitemporelle.

Ce paramètre est facultatif et s’applique uniquement aux tables bitemporelles.

ignore_null_updates

bool

Contrôle la façon dont les valeurs null des mises à jour CDC entrantes sont gérées. Lorsque ignore_null_updates est True, les colonnes null d'une mise à jour entrante sont ignorées ; la valeur existante dans la ligne cible est préservée. Ceci s'applique également aux colonnes imbriquées avec des valeurs null. Lorsque ignore_null_updates est False, les colonnes null d'une mise à jour entrante écrasent les valeurs existantes dans la cible. default to False.

Défini sur True lorsque les événements source n'incluent que les colonnes modifiées, de sorte que les colonnes inchangées ne soient pas écrasées par null.

La default est False.

Pour un contrôle plus précis des colonnes qui ignorent les valeurs null, utilisez ignore_null_updates_column_list ou ignore_null_updates_except_column_list.

ignore_null_updates_column_list

list

Un sous-ensemble de colonnes pour lesquelles les valeurs null dans un enregistrement de modification entrant sont ignorées, de sorte que chacune de ces colonnes conserve sa valeur existante dans la cible. Les colonnes en dehors de la liste appliquent des valeurs null explicites. Utilisez ce paramètre pour appliquer des mises à jour partielles lorsque votre source n'envoie que les colonnes modifiées. Équivalent à la clause SQL IGNORE NULL UPDATES ON columnList. Utilisez soit ignore_null_updates_column_list, soit ignore_null_updates_except_column_list, pas les deux.

ignore_null_updates_except_column_list

list

Un sous-ensemble de colonnes qui appliquent des valeurs null explicites. Toutes les autres colonnes ignorent les valeurs null dans un enregistrement de modification entrant et conservent leur valeur existante dans la cible. Équivalent à la clause SQL IGNORE NULL UPDATES ON * EXCEPT (...). Utilisez ignore_null_updates_column_list ou ignore_null_updates_except_column_list, pas les deux.

columns_to_update

str OU col()

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 (array<string>). Les colonnes ne figurant pas dans le tableau conservent leurs valeurs cibles existantes, tandis que les colonnes répertoriées sont écrites à partir de la source, y compris les valeurs null explicites. Utilisez ce paramètre lorsque chaque enregistrement de modification met à jour un ensemble de colonnes différent et que vous devez appliquer des valeurs null explicites. Équivalent à la clause SQL COLUMNS TO UPDATE.

Vous ne pouvez pas définir columns_to_update avec ignore_null_updates, ignore_null_updates_column_list ou ignore_null_updates_except_column_list. columns_to_update n'est pas pris en charge pour les tables bitemporelles.

apply_as_deletes

str OU expr()

Spécifie quand un événement CDC doit être traité comme un DELETE plutôt qu'un upsert. Vous pouvez spécifier l'une des options suivantes :

  • Une chaîne : "Operation = 'DELETE'"
  • Une fonction expr() de Spark SQL : expr("Operation = 'DELETE'")

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 pipelines.cdc.tombstoneGCThresholdInSeconds.

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 pipelines.cdc.tombstoneGCThresholdInSeconds sur une valeur qui dépasse le délai maximal attendu entre l'arrivée de l'événement et l'exécution du pipeline. Cela garantit que les marqueurs de suppression sont conservés suffisamment longtemps pour gérer correctement les événements de suppression tardifs ou désordonnés.

apply_as_truncates

str OU expr()

Spécifie quand un événement CDC doit être traité comme une table complète TRUNCATE. Vous pouvez spécifier l'une des options suivantes :

  • Une chaîne : "Operation = 'TRUNCATE'"
  • Une fonction expr() de Spark SQL : expr("Operation = 'TRUNCATE'")

É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 apply_as_truncates parameter est pris en charge uniquement pour le type SCD 1. Le type SCD 2 ne prend pas en charge les Opérations de troncation.

column_list OU except_column_list

list

Un sous-ensemble de colonnes à inclure dans la table cible. Utilisez column_list pour spécifier la liste complète des colonnes à inclure. Utilisez except_column_list pour spécifier les colonnes à exclure. Vous pouvez déclarer chaque valeur comme une liste de chaînes de caractères ou comme des fonctions col() Spark SQL :

  • column_list = ["userId", "name", "city"]
  • column_list = [col("userId"), col("name"), col("city")]
  • except_column_list = ["operation", "sequenceNum"]
  • 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 column_list ou except_column_list n'est transmis à la fonction.

stored_as_scd_type

str OU int

Pour stocker des enregistrements en tant que type SCD 1, type SCD 2 ou bitemporels. Définissez sur 1 pour le type SCD 1, 2 pour le type SCD 2, ou "bitemporal" pour suivre les modifications à la fois pour l'heure métier et l'heure système. Le bitemporel nécessite system_sequence_by et est en Bêta. Consultez Bitemporal AUTO CDC. Le default est le type de SCD 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.

name

str

Le nom du flux. Si non spécifié, default est la même valeur que target.

once

bool

Vous pouvez éventuellement définir le flux comme un flux ponctuel, tel qu’un remplissage rétroactif. L'utilisation de once=True modifie le flux de deux manières :

  • La valeur de retour. streaming-query. doit être un DataFrame batch dans ce cas, et non un DataFrame en streaming.
  • Le flux s'exécute une seule fois par default. Si le pipeline est mis à jour avec un refresh complet, alors le flux ONCE s'exécute de nouveau pour recréer les données.
Sur cette page