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 |
INSERT [ONCE] INTO target_table BY NAME query
}

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

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

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

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;