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 |
INSERT [ONCE] INTO target_table BY NAME [ replace_using_spec ] query
}
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.
-
Une instruction
AUTO CDC ... INTOqui définit le flux, avec uncreate_auto_cdc_flow_spec. Vous devez inclure soit une instructionAUTO CDC ... INTO, soit une instructionINSERT INTO. UtilisezAUTO CDC ... INTOlorsque la query source utilise la sémantique de modification des données.Pour plus d’informations, consultez AUTO CDC INTO (pipelines).
-
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
ONCEn'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'optionskipChangeCommitspour gérer les erreurs.INSERT INTOest mutuellement exclusif avecAUTO CDC ... INTO. UtilisezAUTO CDC ... INTOlorsque les données source incluent la fonctionnalité de capture de données modifiées (CDC). UtilisezINSERT INTOlorsque 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
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
ONCEmodifie le flux de deux manières :- La source
queryoucreate_auto_cdc_flow_specn'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
ONCEs'exécute de nouveau pour recréer les données.
ONCEne peut pas être utilisé avecREPLACE USING, qui nécessite une source de streaming. - La source
Exemples
-- 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);