Automatisering van datapijplijnen in Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture

Wat een stream doet

| 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)
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;
SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
METADATA$ACTION: INSERT of DELETEMETADATA$ISUPDATE: TRUE als onderdeel van een update-paarMETADATA$ROW_ID: Unieke fysieke rij-IDMETADATA$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';
Automatisering van datapijplijnen in Snowflake