Aller au contenu principal

AUTO CDC VERS (pipelines)

Utilisez l'instruction AUTO CDC ... INTO pour créer un flux qui utilise la fonctionnalité de capture de données modifiées (CDC) des LakeFlow Pipelines. Cette instruction lit les modifications d'une source CDC et les applique à une cible de streaming.

Syntaxe

CREATE OR REFRESH STREAMING TABLE table_name;

CREATE FLOW flow_name AS AUTO CDC [ONCE] INTO table_name
FROM source
KEYS (keys)
[IGNORE NULL UPDATES [ON {columnList | * EXCEPT (exceptColumnList)}]]
[APPLY AS DELETE WHEN condition]
[APPLY AS TRUNCATE WHEN condition]
SEQUENCE BY orderByColumn
[SYSTEM SEQUENCE BY systemOrderByColumn]
[COLUMNS {columnList | * EXCEPT (exceptColumnList)}]
[STORED AS {SCD TYPE 1 | SCD TYPE 2 | BITEMPORAL}]
[TRACK HISTORY ON {columnList | * EXCEPT (exceptColumnList)}]
[COLUMNS TO UPDATE columnName]

Vous définissez les contraintes de qualité des données pour la cible en utilisant la même clause CONSTRAINT que les autres requêtes de pipeline. Voir Gérer la qualité des données avec les attentes du pipeline.

Le comportement par default pour les événements INSERT et UPDATE est d' upserter les événements CDC 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 la condition APPLY AS DELETE WHEN.

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. Pour les tables SCD de type 2, lors de la spécification du schéma de la table cible, vous devez également inclure les colonnes __START_AT et __END_AT avec le même type de données que le champ sequence_by.

Consultez Les AUTO CDC APIs : simplifiez la capture de données modifiées avec le pipeline.

parameter

  • ONCE

    Spécifier ONCE signifie que cela effectue une insertion unique, ou un remplissage, dans la table cible. Il n’est pas réexécuté si le pipeline dans lequel il se trouve est refresh, sauf en cas de refresh complète.

    Cette clause est facultative.

  • flow_name

    Le nom du flux à créer.

  • source

    La source des données. La source doit être une source de streaming . Utilisez le mot-clé STREAM pour utiliser la sémantique de streaming afin de lire à partir de la source. Si la lecture rencontre une modification ou une suppression d’un enregistrement existant, une erreur est générée. Il est plus sûr de lire à partir de sources statiques ou à ajout seulement. Pour ingérer des données qui contiennent des commits de modification, vous pouvez utiliser Python et l’option skipChangeCommits pour gérer les erreurs.

    Pour plus d'informations sur le streaming de données, consultez Transformer des données avec des pipelines.

  • KEYS

    La colonne ou la combinaison de colonnes qui identifient de manière unique une ligne dans les données source. Les valeurs de ces colonnes sont utilisées pour identifier les événements CDC qui s'appliquent à des enregistrements spécifiques dans la table cible.

    Pour définir une combinaison de colonnes, utilisez une liste de colonnes séparées par des virgules.

    Cette clause est requise.

  • IGNORE NULL UPDATES

    Permet l'ingestion de mises à jour contenant un sous-ensemble des colonnes cibles. Lorsqu'un événement CDC correspond à une ligne existante et que IGNORE NULL UPDATES est spécifié, les colonnes avec une valeur null conservent leurs valeurs existantes dans la cible. Cela s'applique également aux colonnes imbriquées avec une valeur null.

    Pour les mises à jour partielles, ajoutez une clause ON pour contrôler les colonnes qui ignorent les valeurs null :

    • IGNORE NULL UPDATES ON columnList: seules les colonnes répertoriées conservent leurs valeurs existantes lorsque la valeur entrante est null. Toutes les autres colonnes appliquent des valeurs null explicites.
    • IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList): toutes les colonnes, à l'exception de celles répertoriées, conservent leurs valeurs existantes lorsque la valeur entrante est null. Les colonnes listées appliquent des valeurs explicites null.

    Cette clause est facultative.

    La valeur par default est de remplacer les colonnes existantes avec des valeurs null.

  • APPLY AS DELETE WHEN

    Spécifie quand un événement CDC doit être traité comme un DELETE plutôt qu'une upsert.

    Pour les sources de SCD de type 2, afin de gérer les données désordonnées, la ligne supprimée est temporairement conservée en tant que marqueur de suppression dans la table Delta sous-jacente, et une vue est créée dans le metastore qui filtre ces marqueurs de suppression. L'intervalle de rétention peut être configuré avec la propriété de table pipelines.cdc.tombstoneGCThresholdInSeconds.

    Cette clause est facultative.

  • APPLY AS TRUNCATE WHEN

    Spécifie quand un événement CDC doit être traité comme une table complète TRUNCATE. Puisque cette clause déclenche une troncature complète de la table cible, elle ne devrait être utilisée que pour des cas d'utilisation spécifiques nécessitant cette fonctionnalité.

    La clause APPLY AS TRUNCATE WHEN n'est prise en charge que pour le SCD de type 1. Le SCD de type 2 ne prend pas en charge l'opération de troncation.

    Cette clause est facultative.

  • SEQUENCE BY

    Le nom de la colonne spécifiant l'ordre logique des événements CDC dans les données source. Le traitement du pipeline utilise ce séquencement pour gérer les événements de modification qui arrivent dans le désordre.

    Si plusieurs colonnes sont nécessaires pour le séquençage, utilisez une expression STRUCT : elle ordonne d'abord par le premier champ de structure, puis par le second champ en cas d'égalité, et ainsi de suite.

    Les colonnes spécifiées doivent être des types de données triables.

    Cette clause est requise.

  • SYSTEM SEQUENCE BY

info

Bêta

Le CDC AUTO bitemporel est en bêta.

Le nom de la colonne spécifiant l'heure système à laquelle chaque événement CDC est connu du système. Utilisé avec STORED AS BITEMPORAL pour suivre les changements à travers le temps métier (SEQUENCE BY) et l'heure système. Voir Bitemporal AUTO CDC.

Les colonnes spécifiées doivent être des types de données triables.

Cette clause est facultative et ne s'applique qu'aux tables bitemporelles.

  • COLUMNS

    Spécifie un sous-ensemble de colonnes à inclure dans la table cible. Vous pouvez :

    • Spécifiez la liste complète des colonnes à inclure : COLUMNS (userId, name, city).
    • Spécifiez une liste de colonnes à exclure : COLUMNS * EXCEPT (operation, sequenceNum)

    Cette clause est facultative.

    Le default est d'inclure toutes les colonnes dans la table cible lorsque la clause COLUMNS n'est pas spécifiée.

  • STORED AS

    S'il faut stocker les enregistrements en tant que SCD de type 1, SCD de type 2 ou bitemporels.

    Définissez sur BITEMPORAL pour suivre les modifications à la fois au niveau de l'heure métier et de l'heure système. Bitemporal nécessite SYSTEM SEQUENCE BY et est en version bêta. Voir Bitemporal AUTO CDC.

    Cette clause est facultative.

    Le default est le type de SCD 1.

  • TRACK HISTORY ON

    Spécifie un sous-ensemble de colonnes de sortie pour générer des enregistrements d'historique en cas de modifications de ces colonnes spécifiées. Vous pouvez :

    • Spécifiez la liste complète des colonnes à suivre : COLUMNS (userId, name, city).
    • Spécifiez une liste de colonnes à exclure du suivi : COLUMNS * EXCEPT (operation, sequenceNum)

    Cette clause est facultative. Le default est de suivre l'historique pour toutes les colonnes de sortie lorsqu'il y a des changements, équivalent à TRACK HISTORY ON *.

  • COLUMNS TO UPDATE

    Spécifie le nom d'une colonne source qui contient, pour chaque enregistrement de modification, l'ensemble des colonnes à mettre à jour sous la forme d'un tableau de chaînes de noms de colonnes (array<string>). Les colonnes qui ne figurent 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 cette clause pour les mises à jour partielles lorsque chaque enregistrement de modification met à jour un ensemble de colonnes différent et que vous devez appliquer des valeurs null explicites.

    Vous ne pouvez pas utiliser COLUMNS TO UPDATE avec IGNORE NULL UPDATES, et ce n'est pas pris en charge pour les tables bitemporales.

    Cette clause est facultative.

Exemples

SQL
-- Create a streaming table, then use AUTO CDC to populate it:
CREATE OR REFRESH STREAMING TABLE target;

CREATE FLOW flow
AS AUTO CDC INTO
target
FROM stream(cdc_data.users)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2
TRACK HISTORY ON * EXCEPT (city);