Aller au contenu principal

CRÉER UN FLUX (pipelines)

Utilisez l'instruction CREATE FLOW pour créer des flux ou des remplissages pour les tables d'un pipeline.

Syntaxe​

CREATE FLOW flow_name [COMMENT comment] AS
{
AUTO CDC [ONCE] INTO target_table create_auto_cdc_flow_spec |
AUTO CDC [ONCE] INTO target_table create_auto_cdc_from_snapshot_spec |
INSERT [ONCE] INTO target_table BY NAME [ replace_using_spec ] query
}

create_auto_cdc_from_snapshot_spec
FROM SNAPSHOT ( snapshot_query )
[ WITH VERSION ( version_query ) ]
KEYS ( key [, ...] )
[ STORED AS { SCD TYPE 1 | SCD TYPE 2 } ]
[ TRACK HISTORY ON { col_list | * EXCEPT ( col_list ) } ]

replace_using_spec
REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column

parameter​

  • flow_name

    Le nom du flux à créer.

  • Commentaire

    Une description facultative pour le flux.

  • AUTO CDC INTO

    Une instruction AUTO CDC ... INTO qui définit le flux, avec un create_auto_cdc_flow_spec. Vous devez inclure soit une instruction AUTO CDC ... INTO, soit une instruction INSERT INTO. Utilisez AUTO CDC ... INTO lorsque la query source utilise la sémantique de modification des données.

    Pour plus d’informations, consultez AUTO CDC INTO (pipelines).

  • AUTO CDC ... FROM SNAPSHOT

    Une instruction AUTO CDC ... INTO qui dérive des changements en comparant des instantanés au lieu de lire un flux de modifications. Utilisez ce formulaire lorsque la capture de données modificatives (CDC) n'est pas activée sur la source et que seuls des instantanés complets sont disponibles. La source est spécifiée en deux parties : une clause obligatoire FROM SNAPSHOT (snapshot_query) qui lit les données de l'instantané, et une clause facultative WITH VERSION (version_query) qui sélectionne la version d'instantané suivante à traiter. Voir Fonctionnement de AUTO CDC FROM SNAPSHOT.

    • FROM SNAPSHOT (snapshot_query)

      Obligatoire. Une query qui lit les données d’instantané pour la version sélectionnée par WITH VERSION (...). Le moteur compare le résultat à l’instantané précédemment validé pour en dériver les insertions, les mises à jour et les suppressions, puis les fusionne dans la cible en utilisant KEYS pour l’identité des lignes et STORED AS pour déterminer la manière dont les modifications sont stockées.

      Appelez current_snapshot_version() dans cette query pour faire référence à la version sélectionnée par WITH VERSION (...). Si WITH VERSION (...) n’est pas spécifié, current_snapshot_version() n’est pas appelable dans FROM SNAPSHOT (...).

      Lorsque WITH VERSION (...) est omis, le moteur lit la source directement par l’intermédiaire de FROM SNAPSHOT (...), et la query de l’instantané s’exécute uniquement lors du chargement initial, tandis que la cible ne comporte aucune donnée validée ni aucun état d’instantané validé. Lors de toute mise à jour ultérieure, lorsque la cible contient déjà des données ou a validé un état d’instantané, le flux échoue avec AUTO_CDC_FROM_SNAPSHOT_NON_EMPTY_TARGET_WITHOUT_VERSION. Pour traiter les instantanés sur plusieurs mises à jour, utilisez WITH VERSION (...).

    • WITH VERSION (version_query)

      Facultatif. Une query qui sélectionne la prochaine version d'instantané à traiter. Elle doit retourner exactement une colonne d'un type pouvant être ordonné et 0 ou 1 ligne. Lorsqu'elle retourne 1 ligne, la valeur ne doit pas être nulle. La colonne peut être une valeur scalaire, telle qu'un BIGINT, ou un STRUCT dont tous les champs peuvent être ordonnés. Une query de version qui retourne plus d'une colonne, plus d'une ligne ou une valeur nulle fait échouer le flux avec INVALID_AUTO_CDC_FROM_SNAPSHOT_VERSION_QUERY.

      Au cours d’une mise à jour de pipeline, le moteur répète les étapes suivantes : il évalue la query de version ; si la query renvoie 0 ligne, il arrête de traiter ce flux pour la mise à jour actuelle ; si la query renvoie 1 ligne, le moteur expose cette valeur via current_snapshot_version(), évalue la query d’instantané, commit l’instantané résultant et expose la version validée via last_snapshot_version(). Le moteur réévalue ensuite la query de version pour sélectionner la version suivante. Une seule mise à jour du pipeline traite les versions dans l’ordre jusqu’à ce que la query de version ne renvoie aucune ligne.

      Chaque version retournée après un commit réussi doit être supérieure à la version précédemment validée ; une version qui n'augmente pas fait échouer la mise à jour avec APPLY_CHANGES_FROM_SNAPSHOT_ERROR.OUT_OF_ORDER_SNAPSHOT_VERSION. Le type de données de la valeur de version doit rester inchangé entre les commits d'instantané ; une modification du type de données fait échouer la mise à jour avec AUTO_CDC_FROM_SNAPSHOT_VERSION_SCHEMA_CHANGED. Un full refresh efface l'état de version persisté.

    • KEYS

      Obligatoire. Les colonnes de clé primaire utilisées pour identifier les lignes entre les instantanés pour la détection des changements.

    • STORED AS { SCD TYPE 1 | SCD TYPE 2 }

      Facultatif. Spécifie la manière dont les modifications sont stockées dans la table cible. La valeur default est SCD TYPE 1.

    • TRACK HISTORY ON { col_list | * EXCEPT (col_list) }

      Facultatif. Ne s’applique qu’avec SCD TYPE 2. Spécifie les colonnes qui Trigger une nouvelle ligne d’historique lorsqu’elles changent. Fournissez soit une liste explicite de colonnes, soit * EXCEPT (col_list) pour suivre chaque colonne à l’exception de celles répertoriées.

    La CDC par instantané ne prend pas en charge WHERE ou SEQUENCE BY. L'ordre inter-instantanés est exprimé via WITH VERSION (...).

  • table cible

    La table à mettre à jour. Il doit s'agir d'une table de streaming.

  • INSERT INTO

    Définit une query de table qui est insérée dans la table cible. Si l'option ONCE n'est pas fournie, la query doit être une query 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 ont des commits de modifications, vous pouvez utiliser Python et l'option skipChangeCommits pour gérer les erreurs.

    INSERT INTO est mutuellement exclusif avec AUTO CDC ... INTO. Utilisez AUTO CDC ... INTO lorsque les données source incluent la fonctionnalité de capture de données modifiées (CDC). Utilisez INSERT INTO lorsque la source ne le contient pas.

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

  • REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column

info

Bêta

Cette fonctionnalité est en version bêta. Nécessite Databricks Runtime 18.2 et versions ultérieures.

Définit le flux comme un flux REPLACE USING, qui remplace toutes les lignes de la table cible correspondant aux colonnes de clé spécifiées et laisse toutes les autres lignes intactes. Utilisez REPLACE USING lorsque votre source est une série de snapshots partiels indexés par colonne. SEQUENCE BY ordonne les mises à jour de sorte que la séquence la plus élevée pour une clé l'emporte, même lorsque les mises à jour arrivent dans le désordre.

Spécifiez au moins une colonne clé et exactement une colonne SEQUENCE BY. La query doit être une query de streaming, et BY NAME est requis. REPLACE USING ne peut pas être combiné avec ONCE ou avec AUTO CDC ... INTO.

Pour plus d’informations, consultez Remplacement de snapshot partiel avec les flux REPLACE USING.

  • Une fois

    Définissez éventuellement le flux comme un flux unique, tel qu'un remplissage. L'utilisation de ONCE modifie le flux de deux manières :

    • La source query ou create_auto_cdc_flow_spec n'est pas une table de 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.

    ONCE ne peut pas être utilisé avec REPLACE USING, qui nécessite une source de streaming.

Exemples​

SQL
-- EXAMPLE 1:
-- Create a streaming table, and add two flows that append data to it:
CREATE OR REFRESH STREAMING TABLE users;

-- first flow into target_table:
CREATE FLOW users_flow AS
INSERT INTO users BY NAME
SELECT * FROM stream(raw_data.users);

-- second flow into target_table:
CREATE FLOW backfill_users AS
INSERT ONCE INTO users BY NAME
SELECT * FROM user_backfill_table;

-- EXAMPLE 2:
-- Create a streaming table, and add a flow that applies CDC changes to it:
CREATE OR REFRESH STREAMING TABLE admins_cdc_target_table;

-- first flow into target_table:
CREATE FLOW admin_cdc_flow AS
AUTO CDC INTO admins_cdc_target_table
FROM stream(cdc_data.admins)
KEYS (userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;

-- EXAMPLE 3:
-- Create a streaming table, and add a REPLACE USING flow that keeps the latest
-- row for each payment_id from a stream of partial snapshots:
CREATE OR REFRESH STREAMING TABLE payments_latest;

CREATE FLOW payments_replace_flow AS
INSERT INTO payments_latest BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

-- EXAMPLE 4:
-- AUTO CDC FROM SNAPSHOT without WITH VERSION: a one-time initial load from a snapshot table.
-- To process later snapshots on each update, add WITH VERSION (see EXAMPLE 5).
CREATE STREAMING TABLE users (user_id INT, name STRING, email STRING);

CREATE FLOW users_snapshot_flow AS
AUTO CDC ONCE INTO users
FROM SNAPSHOT (SELECT * FROM catalog.schema.users_snapshot)
KEYS (user_id)
STORED AS SCD TYPE 1;

-- EXAMPLE 5:
-- AUTO CDC FROM SNAPSHOT with WITH VERSION: pick the next file, then read it as the snapshot:
CREATE STREAMING TABLE orders (order_id INT, product STRING, quantity INT, order_date DATE);

CREATE FLOW orders_cdc AS
AUTO CDC INTO orders
FROM SNAPSHOT (
SELECT order_id, product, quantity, order_date
FROM read_files('/Volumes/catalog/schema/landing/orders/', format => 'json')
WHERE _metadata.file_path = (SELECT version.path FROM current_snapshot_version())
)
WITH VERSION (
SELECT struct(modification_time, path) AS version
FROM list_files('/Volumes/catalog/schema/landing/orders/')
WHERE (
NOT EXISTS (SELECT 1 FROM last_snapshot_version())
OR struct(modification_time, path) > (SELECT version FROM last_snapshot_version())
)
ORDER BY modification_time, path
LIMIT 1
)
KEYS (order_id)
STORED AS SCD TYPE 2;

-- EXAMPLE 6:
-- One-time snapshot backfill plus a streaming CDC flow into the same target.
-- The backfill omits WITH VERSION, so it uses an implicit timestamp version. The
-- streaming flow's SEQUENCE BY column (event_ts) must be a TIMESTAMP so its type
-- matches that implicit version on the shared target.
CREATE STREAMING TABLE customers (
customer_id INT, name STRING, email STRING, address STRING, event_ts TIMESTAMP
);

CREATE FLOW customers_snapshot_backfill AS
AUTO CDC ONCE INTO customers
FROM SNAPSHOT (SELECT * FROM catalog.schema.customers_snapshot)
KEYS (customer_id)
STORED AS SCD TYPE 1;

CREATE FLOW customers_cdc AS
AUTO CDC INTO customers
FROM STREAM(customers_cdc_events)
KEYS (customer_id)
SEQUENCE BY event_ts
STORED AS SCD TYPE 1;

Pour en savoir plus sur l'association d'un remplissage ponctuel et de la CDC continue sur la même cible, consultez Remplissage de données historiques avec des pipelines.