Автоматизация конвейеров данных в Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: захват изменений данных

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

| Тип потока | Захватывает | Лучше использовать для |
|---|---|---|
| Standard | Все типы таблиц и представления, все DML-изменения — INSERT, UPDATE, DELETE | Таблицы с изменяемыми строками (например, отгрузки) |
| Append-only | Все типы таблиц и представления, кроме внешних — только INSERT | Таблицы с однократной записью (например, события доставки) — эффективнее |
| Insert-only | Внешние таблицы Apache Iceberg — только INSERT | Внешние таблицы |
Таблицы каталогов предоставляют метаданные файлов стейджа (имя, размер, дата изменения)
Стандартный поток на таблице отгрузок
CREATE STREAM shipments_stream
ON TABLE logistics.shipments;
Поток append-only на таблице событий доставки
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 или DELETEMETADATA$ISUPDATE: TRUE при наличии пары обновленийMETADATA$ROW_ID: Уникальный идентификатор строкиMETADATA$ISUPDATE = TRUE


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