Streamy a Change Data Capture

Automatizace datových pipeline ve Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Zpracování pouze změněných dat

CDC: Change Data Capture Srovnání dvou způsobů zpracování dat

Automatizace datových pipeline ve Snowflake

Streamy

Co stream dělá

  • Sleduje každý INSERT, UPDATE a DELETE ve zdrojové tabulce

Screenshot 2026-05-11 at 10.50.31 am.png

  • Udržuje průběžný protokol změn — bez duplikace dat
  • Po zpracování se offset posune; další čtení začíná od začátku
1 * Snowflake Learning Resource
Automatizace datových pipeline ve Snowflake

Typy streamů

 

Typ streamu Zachycuje Vhodné pro
Standard Všechny typy tabulek a pohledy & všechny změny DML – sleduje vkládání, aktualizace, mazání Tabulky, kde se může změnit libovolný řádek (např. zásilky)
Append-only Všechny typy tabulek a pohledy kromě externích – sleduje pouze vkládání řádků Tabulky pouze pro vkládání (např. události doručení) – efektivnější
Insert-only Externě spravované tabulky Apache Iceberg a externí tabulky – sleduje pouze vkládání řádků Externí tabulky

Adresářové tabulky poskytují metadata souborů pro stage (název, velikost, čas poslední úpravy)

Automatizace datových pipeline ve Snowflake

Vytvoření streamu

Standardní stream na tabulce zásilek

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Stream append-only na tabulce událostí doručení

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Automatizace datových pipeline ve Snowflake

Sloupce metadat streamu

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT nebo DELETE
  • METADATA$ISUPDATE: TRUE, pokud je součástí páru aktualizace
  • METADATA$ROW_ID: Jedinečný fyzický identifikátor řádku
  • Aktualizace se zobrazují jako pár DELETE + INSERT, oba označené METADATA$ISUPDATE = TRUE

Sloupce metadat streamu zobrazující METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID s ukázkovými daty

Automatizace datových pipeline ve Snowflake

Offset streamu

Diagram časové osy – Stream vytvořen (offset začíná zde) → Změny probíhají ve zdrojové tabulce (stream shromažďuje záznamy) → Stream je zpracován v transakci (offset se posouvá na aktuální čas)

Automatizace datových pipeline ve Snowflake

Streamy v pipeline – přehled

  • Streamy se kombinují s tasky — objekty Snowflake, které spouštějí SQL podle plánu
  • Task čte pouze změněné řádky; při 10 mil. řádcích: 2 sekundy vs. 2 minuty

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Automatizace datových pipeline ve Snowflake

Streamy v pipeline – dotaz

  • Streamy se kombinují s tasky – objekty Snowflake, které spouštějí SQL podle plánu
  • Task čte pouze změněné řádky; při 10 mil. řádcích: 2 sekundy vs. 2 minuty
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';
Automatizace datových pipeline ve Snowflake

Pojďme si procvičit!

Automatizace datových pipeline ve Snowflake

Preparing Video For Download...