Strömmar och Change Data Capture

Automatisering av datapipelines i Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Bearbeta bara det som ändrats

CDC: Change Data Capture Två jämförelser av hur data bearbetas

Automatisering av datapipelines i Snowflake

Strömmar

Vad en ström gör

  • Spårar varje INSERT, UPDATE och DELETE i en källtabell

Screenshot 2026-05-11 at 10.50.31 am.png

  • Upprätthåller en löpande ändringslogg — ingen dataduplicering
  • När strömmen konsumeras avancerar offset; nästa läsning börjar om
1 * Snowflake Learning Resource
Automatisering av datapipelines i Snowflake

Strömstyper

 

Strömstyp Fångar Bäst för
Standard Alla tabelltyper och vyer samt alla DML-ändringar – spårar insättningar, uppdateringar och borttagningar Tabeller där en rad kan ändras (t.ex. leveranser)
Append-only Alla tabelltyper och vyer utom externa tabeller – spårar bara radinsättningar Tabeller med engångsinlägg (t.ex. leveranshändelser) – mer effektivt
Insert-only Externt hanterade Apache Iceberg- och externa tabeller – spårar bara radinsättningar Externa tabeller

Katalogtabeller visar filmetadata för en stage (namn, storlek, senast ändrad)

Automatisering av datapipelines i Snowflake

Skapa en ström

Standardström på leveranstabellen

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Append-only-ström på leveranshändelsetabellen

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Automatisering av datapipelines i Snowflake

Strömmens metadatakolumner

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT eller DELETE
  • METADATA$ISUPDATE: TRUE när den ingår i ett uppdateringspar
  • METADATA$ROW_ID: Unik fysisk radidentifierare
  • Uppdateringar visas som ett DELETE + INSERT-par, båda flaggade med METADATA$ISUPDATE = TRUE

Strömmens metadatakolumner med METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID och exempeldata

Automatisering av datapipelines i Snowflake

Strömmens offset

Tidslinjediagram – Ström skapas (offset börjar här) → Ändringar sker i källtabellen (strömmen samlar poster) → Strömmen konsumeras i en transaktion (offset avancerar till nu)

Automatisering av datapipelines i Snowflake

Strömmar i en pipeline – översikt

  • Strömmar kombineras med tasks — Snowflake-objekt som kör SQL enligt schema
  • Task läser bara ändrade rader; med 10 M rader: 2 sekunder vs 2 minuter

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Automatisering av datapipelines i Snowflake

Strömmar i en pipeline – fråga

  • Strömmar kombineras med tasks – Snowflake-objekt som kör SQL enligt schema
  • Task läser bara ändrade rader; med 10 M rader: 2 sekunder vs 2 minuter
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 av datapipelines i Snowflake

Lass uns üben!

Automatisering av datapipelines i Snowflake

Preparing Video For Download...