Потоки и захват изменений данных

Автоматизация конвейеров данных в Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Обработка только изменённых данных

CDC: захват изменений данных Два способа обработки данных

Автоматизация конвейеров данных в Snowflake

Потоки

Что делает поток

  • Отслеживает каждый INSERT, UPDATE и DELETE в исходной таблице

Screenshot 2026-05-11 at 10.50.31 am.png

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

Типы потоков

 

Тип потока Захватывает Лучше использовать для
Standard Все типы таблиц и представления, все DML-изменения — INSERT, UPDATE, DELETE Таблицы с изменяемыми строками (например, отгрузки)
Append-only Все типы таблиц и представления, кроме внешних — только INSERT Таблицы с однократной записью (например, события доставки) — эффективнее
Insert-only Внешние таблицы Apache Iceberg — только INSERT Внешние таблицы

Таблицы каталогов предоставляют метаданные файлов стейджа (имя, размер, дата изменения)

Автоматизация конвейеров данных в Snowflake

Создание потока

Стандартный поток на таблице отгрузок

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Поток append-only на таблице событий доставки

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

Смещение потока

Диаграмма: Поток создан (смещение начинается здесь) → Изменения в исходной таблице (поток накапливает записи) → Поток потреблён в транзакции (смещение продвигается вперёд)

Автоматизация конвейеров данных в Snowflake

Потоки в конвейере: обзор

  • Потоки используются совместно с задачами — объектами Snowflake, выполняющими SQL по расписанию
  • Задача читает только изменённые строки; при 10 млн строк: 2 секунды вместо 2 минут

Screenshot 2026-05-11 at 10.48.51 am.png

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

Потоки в конвейере: запрос

  • Потоки используются совместно с задачами — объектами Snowflake, выполняющими SQL по расписанию
  • Задача читает только изменённые строки; при 10 млн строк: 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

Lass uns üben!

Автоматизация конвейеров данных в Snowflake

Preparing Video For Download...