Streams und Change Data Capture

Automatisierung von Datenpipelines in Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Nur Geändertes verarbeiten

CDC: Change Data Capture Zwei Vergleiche, wie Daten verarbeitet werden

Automatisierung von Datenpipelines in Snowflake

Streams

Was ein Stream macht

  • Protokolliert jeden INSERT, UPDATE und DELETE auf einer Quelltabelle

Screenshot 2026-05-11 at 10.50.31 am.png

  • Führt ein fortlaufendes Änderungsprotokoll — keine Datenkopie
  • Nach dem Konsumieren springt der Offset vor; nächster Read startet frisch
1 * Snowflake Learning Resource
Automatisierung von Datenpipelines in Snowflake

Stream-Typen

 

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)

Automatisierung von Datenpipelines in Snowflake

Einen Stream erstellen

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

Stream-Metadaten-Spalten

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT oder DELETE
  • METADATA$ISUPDATE: TRUE, wenn Teil eines Update-Paars
  • METADATA$ROW_ID: Eindeutige physische Zeilen-ID
  • Updates erscheinen als DELETE + INSERT-Paar, beide mit METADATA$ISUPDATE = TRUE markiert

Stream-Metadaten-Spalten mit METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID und Beispieldaten

Automatisierung von Datenpipelines in Snowflake

Der Stream-Offset

Zeitachse – Stream erstellt (Offset startet hier) → Änderungen in der Quelltabelle (Stream sammelt Datensätze) → Stream in einer Transaktion konsumiert (Offset springt auf jetzt)

Automatisierung von Datenpipelines in Snowflake

Streams im Pipeline-Überblick

  • Streams arbeiten mit Tasks — Snowflake-Objekte, die SQL nach Zeitplan ausführen
  • Task liest nur geänderte Zeilen; bei 10 Mio. Zeilen: 2 Sekunden statt 2 Minuten

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Automatisierung von Datenpipelines in Snowflake

Stream in einer Pipeline: Query

  • Streams arbeiten mit Tasks – Snowflake-Objekte, die SQL nach Zeitplan ausführen
  • Task liest nur geänderte Zeilen; bei 10 Mio. Zeilen: 2 Sekunden statt 2 Minuten
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

Lass uns üben!

Automatisierung von Datenpipelines in Snowflake

Preparing Video For Download...