Flux et Change Data Capture

Automatisation des pipelines de données dans Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Traiter uniquement ce qui a changé

CDC : Change Data Capture Deux comparaisons sur la façon d’exécuter les données

Automatisation des pipelines de données dans Snowflake

Flux

Ce que fait un flux

  • Suit chaque INSERT, UPDATE et DELETE sur une table source

Capture d’écran 2026-05-11 à 10.50.31.png

  • Tient un journal des changements — sans dupliquer les données
  • Une fois consommé, le décalage avance ; la lecture suivante repart à zéro
1 * Snowflake Learning Resource
Automatisation des pipelines de données dans Snowflake

Types de flux

 

Type de flux Capture Idéal pour
Standard Tous types de tables et vues & tous changements DML - insertions, mises à jour, suppressions Tables où toute ligne peut changer (ex. expéditions)
Ajout-seulement Tous types de tables et vues, sauf tables externes - n’enregistre que les insertions Tables à insertion unique (ex. événements de livraison) - plus efficace
Insertion-seulement Tables Apache Iceberg gérées externement et tables externes - n’enregistre que les insertions Tables externes

Les tables d’annuaire exposent les métadonnées de fichiers d’un stage (nom, taille, horodatage de dernière modification)

Automatisation des pipelines de données dans Snowflake

Créer un flux

Flux standard sur la table shipments

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Flux ajout-seulement sur la table delivery_events

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Automatisation des pipelines de données dans Snowflake

Colonnes de métadonnées du flux

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION : INSERT ou DELETE
  • METADATA$ISUPDATE : TRUE quand c’est une paire de mise à jour
  • METADATA$ROW_ID : Identifiant physique unique de ligne
  • Les mises à jour apparaissent en DELETE + INSERT, les deux avec METADATA$ISUPDATE = TRUE

Colonnes de métadonnées de flux affichant METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID avec des exemples

Automatisation des pipelines de données dans Snowflake

Le décalage (offset) du flux

Schéma temporel - Flux créé (le décalage commence ici) → Changements dans la table source (le flux accumule) → Flux consommé dans une transaction (le décalage avance jusqu’à maintenant)

Automatisation des pipelines de données dans Snowflake

Aperçu : flux dans un pipeline

  • Les flux fonctionnent avec les tâches — objets Snowflake qui exécutent du SQL planifié
  • La tâche ne lit que les lignes modifiées ; avec 10 M de lignes : 2 s vs 2 min

Capture d’écran 2026-05-11 à 10.48.51.png

1 * Snowflake Learning Resource
Automatisation des pipelines de données dans Snowflake

Requête : flux dans un pipeline

  • Les flux fonctionnent avec les tâches - objets Snowflake qui exécutent du SQL planifié
  • La tâche ne lit que les lignes modifiées ; avec 10 M de lignes : 2 s vs 2 min
CREATE TASK logistics.sync_shipments
  WAREHOUSE = compute_wh
  SCHEDULE = '5 MINUTE'
  WHEN SYSTEM$STREAM_HAS_DATA('logistics.staging_shipments_stream')
AS
  INSERT INTO logistics.shipments
  SELECT shipment_id, region, carrier, delivery_days
  FROM logistics.staging_shipments_stream
  WHERE METADATA$ACTION = 'INSERT';
Automatisation des pipelines de données dans Snowflake

Passons à la pratique !

Automatisation des pipelines de données dans Snowflake

Preparing Video For Download...