Strumienie i przechwytywanie zmian danych

Automatyzacja potoków danych w Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Przetwarzanie tylko zmienionych danych

CDC: Change Data Capture Dwa porównania sposobu przetwarzania danych

Automatyzacja potoków danych w Snowflake

Strumienie

Działanie strumienia

  • Śledzi każde INSERT, UPDATE i DELETE w tabeli źródłowej

Screenshot 2026-05-11 at 10.50.31 am.png

  • Prowadzi bieżący dziennik zmian — bez duplikowania danych
  • Po skonsumowaniu offset przesuwa się; kolejny odczyt zaczyna od nowa
1 * Snowflake Learning Resource
Automatyzacja potoków danych w Snowflake

Typy strumieni

 

Typ strumienia Przechwytuje Najlepsze zastosowanie
Standard Wszystkie typy tabel i widoki oraz wszystkie zmiany DML – śledzi INSERT, UPDATE, DELETE Tabele, w których każdy wiersz może ulec zmianie (np. przesyłki)
Append-only Wszystkie typy tabel i widoki, z wyjątkiem tabel zewnętrznych – śledzi tylko INSERT Tabele z jednorazowym wstawianiem (np. zdarzenia dostawy) – wydajniejsze
Insert-only Zewnętrzne tabele Apache Iceberg i tabele zewnętrzne – śledzi tylko INSERT Tabele zewnętrzne

Tabele katalogów udostępniają metadane plików dla stage'a (nazwa, rozmiar, data modyfikacji)

Automatyzacja potoków danych w Snowflake

Tworzenie strumienia

Strumień standardowy na tabeli shipments

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Strumień append-only na tabeli zdarzeń dostawy

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Automatyzacja potoków danych w Snowflake

Kolumny metadanych strumienia

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT lub DELETE
  • METADATA$ISUPDATE: TRUE jeśli część pary aktualizacji
  • METADATA$ROW_ID: Unikalny fizyczny identyfikator wiersza
  • Aktualizacje są widoczne jako para DELETE + INSERT, obie z flagą METADATA$ISUPDATE = TRUE

Kolumny metadanych strumienia: METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID z przykładowymi danymi

Automatyzacja potoków danych w Snowflake

Offset strumienia

Diagram osi czasu – Utworzenie strumienia (offset startuje tutaj) → Zmiany zachodzą w tabeli źródłowej (strumień gromadzi rekordy) → Strumień skonsumowany w transakcji (offset przesuwa się do teraz)

Automatyzacja potoków danych w Snowflake

Strumienie w potoku – przegląd

  • Strumienie współpracują z zadaniami — obiektami Snowflake uruchamiającymi SQL według harmonogramu
  • Zadanie odczytuje tylko zmienione wiersze; przy 10M wierszy: 2 sekundy zamiast 2 minut

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Automatyzacja potoków danych w Snowflake

Strumienie w potoku – zapytanie

  • Strumienie współpracują z zadaniami - obiektami Snowflake uruchamiającymi SQL według harmonogramu
  • Zadanie odczytuje tylko zmienione wiersze; przy 10M wierszy: 2 sekundy zamiast 2 minut
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';
Automatyzacja potoków danych w Snowflake

Czas na ćwiczenia!

Automatyzacja potoków danych w Snowflake

Preparing Video For Download...