Streams en Change Data Capture

Automatisering van datapijplijnen in Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Alleen wijzigingen verwerken

CDC: Change Data Capture Twee manieren om data te verwerken

Automatisering van datapijplijnen in Snowflake

Streams

Wat een stream doet

  • Houdt elke INSERT, UPDATE en DELETE op een brontabel bij

Screenshot 2026-05-11 at 10.50.31 am.png

  • Beheert een doorlopend changelog — geen dataduplicatie
  • Na verbruik schuift de offset door; volgende read start vers
1 * Snowflake Learning Resource
Automatisering van datapijplijnen in Snowflake

Streamtypes

 

Streamtype Legt vast Beste voor
Standaard Alle tabeltypen en views & alle DML-wijzigingen - volgt inserts, updates, deletes Tabellen waar elke rij kan wijzigen (bijv. shipments)
Append-only Alle tabeltypen en views, behalve external tables - volgt alleen rij-inserts Insert-once-tabellen (bijv. delivery events) - efficiënter
Insert-only Extern beheerde Apache Iceberg- en external tables - volgt alleen rij-inserts Externe tabellen

Directory tables tonen bestandsmetadata voor een stage (naam, grootte, laatst gewijzigd)

Automatisering van datapijplijnen in Snowflake

Een stream maken

Standaardstream op de shipments-tabel

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Append-only-stream op de delivery_events-tabel

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Automatisering van datapijplijnen in Snowflake

Stream-metadatakolommen

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT of DELETE
  • METADATA$ISUPDATE: TRUE als onderdeel van een update-paar
  • METADATA$ROW_ID: Unieke fysieke rij-ID
  • Updates verschijnen als een DELETE + INSERT-paar, beide met METADATA$ISUPDATE = TRUE

Stream-metadatakolommen met METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID en voorbeelddata

Automatisering van datapijplijnen in Snowflake

De stream-offset

Tijdlijnschema - Stream gemaakt (offset start hier) → Wijzigingen in brontabel (stream verzamelt records) → Stream verbruikt in transactie (offset springt naar nu)

Automatisering van datapijplijnen in Snowflake

Streams in een pipeline: overzicht

  • Streams werken samen met tasks — Snowflake-objecten die SQL gepland draaien
  • Task leest alleen gewijzigde rijen; bij 10M rijen: 2 seconden vs 2 minuten

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Automatisering van datapijplijnen in Snowflake

Streams in een pipeline: query

  • Streams werken samen met tasks - Snowflake-objecten die SQL gepland draaien
  • Task leest alleen gewijzigde rijen; bij 10M rijen: 2 seconden vs 2 minuten
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';
Automatisering van datapijplijnen in Snowflake

Laten we oefenen!

Automatisering van datapijplijnen in Snowflake

Preparing Video For Download...