Потоки та Change Data Capture

Автоматизація конвеєрів даних у Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Обробляємо лише зміни

CDC: Change Data Capture Дві схеми порівняння способів обробки даних

Автоматизація конвеєрів даних у Snowflake

Потоки

Що робить потік

  • Відстежує кожен INSERT, UPDATE і DELETE у вихідній таблиці

Screenshot 2026-05-11 at 10.50.31 am.png

  • Веде поточний журнал змін — без дублювання даних
  • Після зчитування зсув просувається; наступне читання починається з нуля
1 * Snowflake Learning Resource
Автоматизація конвеєрів даних у Snowflake

Типи потоків

 

Тип потоку Фіксує Найкраще для
Standard Усі типи таблиць і подань та всі DML-зміни — вставки, оновлення, видалення Таблиці, де може змінитися будь-який рядок (напр. shipments)
Append-only Усі типи таблиць і подань, крім external tables — лише вставки рядків Таблиці «вставка-один-раз» (напр. delivery events) — ефективніше
Insert-only Зовнішньо керовані Apache Iceberg і external tables — лише вставки рядків Зовнішні таблиці

Directory tables показують метадані файлів для сховища (назва, розмір, мітка часу останньої зміни)

Автоматизація конвеєрів даних у Snowflake

Створення потоку

Стандартний потік для таблиці shipments

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Append-only потік для таблиці delivery events

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Автоматизація конвеєрів даних у Snowflake

Стовпці метаданих потоку

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT або DELETE
  • METADATA$ISUPDATE: TRUE, якщо частина пари оновлення
  • METADATA$ROW_ID: Унікальний фізичний ідентифікатор рядка
  • Оновлення відображаються як пара DELETE + INSERT, обидва з METADATA$ISUPDATE = TRUE

Стовпці метаданих потоку: METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID із прикладом даних

Автоматизація конвеєрів даних у Snowflake

Зсув потоку (Stream Offset)

Діаграма часової шкали — Створено потік (відлік починається тут) → Зміни у вихідній таблиці (потік накопичує записи) → Потік спожито в транзакції (зсув просувається до зараз)

Автоматизація конвеєрів даних у Snowflake

Огляд: потоки в конвеєрі

  • Потоки поєднуються із tasks — об’єктами Snowflake, що запускають SQL за розкладом
  • Task читає лише змінені рядки; для 10M рядків: 2 секунди проти 2 хвилин

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Автоматизація конвеєрів даних у Snowflake

Запит: потоки в конвеєрі

  • Потоки поєднуються з tasks — об’єктами Snowflake для запуску SQL за розкладом
  • Task читає лише змінені рядки; для 10M рядків: 2 с проти 2 хв
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';
Автоматизація конвеєрів даних у Snowflake

¡Vamos a practicar!

Автоматизація конвеєрів даних у Snowflake

Preparing Video For Download...