Flux et capture de données modifiées (CDC)

Automatisation des pipelines de données dans Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Traiter seulement ce qui a changé

CDC : capture de données modifiées Deux comparaisons sur la façon d’exécuter les données

Automatisation des pipelines de données dans Snowflake

Flux (Streams)

Ce que fait un flux

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

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

  • Tient un journal des changements — pas de duplication de données
  • Une fois consommé, le décalage avance; la prochaine lecture repart à neuf
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ù n’importe quelle ligne peut changer (ex. expéditions)
Ajouts seulement Tous types de tables et vues, sauf tables externes — insertions uniquement Tables en insertion unique (ex. événements de livraison) — plus efficace
Insert-only Tables Apache Iceberg gérées à l’externe et tables externes — insertions uniquement Tables externes

Les tables d’annuaire exposent les métadonnées de fichiers d’un stage (nom, taille, horodatage de 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 ajouts 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 d’un 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 la ligne
  • Les mises à jour apparaissent en DELETE + INSERT, toutes deux avec METADATA$ISUPDATE = TRUE

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

Automatisation des pipelines de données dans Snowflake

Le décalage (offset) d’un flux

Diagramme chronologique - 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 s’associent aux tâches — objets Snowflake qui exécutent du SQL selon un horaire
  • La tâche lit seulement les lignes modifiées; pour 10 M de lignes : 2 s vs 2 min

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

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

Requête : flux dans un pipeline

  • Les flux s’associent aux tâches — objets Snowflake qui exécutent du SQL selon un horaire
  • La tâche lit seulement les lignes modifiées; pour 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...