Automatizace datových pipeline ve Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture

Co stream dělá

| 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)
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;
SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
METADATA$ACTION: INSERT nebo DELETEMETADATA$ISUPDATE: TRUE, pokud je součástí páru aktualizaceMETADATA$ROW_ID: Jedinečný fyzický identifikátor řádkuMETADATA$ISUPDATE = TRUE


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