Automatisierung von Datenpipelines in Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture

Was ein Stream macht

| Stream-Typ | Erfasst | Geeignet für |
|---|---|---|
| Standard | Alle Tabellentypen und Views & alle DML-Änderungen – erfasst Inserts, Updates, Deletes | Tabellen, in denen sich beliebige Zeilen ändern (z. B. Sendungen) |
| Append-only | Alle Tabellentypen und Views, außer External Tables – erfasst nur Inserts | Insert-once-Tabellen (z. B. Lieferereignisse) – effizienter |
| Insert-only | Extern verwaltete Apache-Iceberg- und External Tables – erfasst nur Inserts | External Tables |
Directory Tables stellen Dateimetadaten für eine Stage bereit (Name, Größe, letztes Änderungsdatum)
Standard-Stream für die Tabelle shipments
CREATE STREAM shipments_stream
ON TABLE logistics.shipments;
Append-only-Stream für die Tabelle delivery_events
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 oder DELETEMETADATA$ISUPDATE: TRUE, wenn Teil eines Update-PaarsMETADATA$ROW_ID: Eindeutige physische Zeilen-IDMETADATA$ISUPDATE = TRUE markiert


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';
Automatisierung von Datenpipelines in Snowflake