Автоматизація конвеєрів даних у Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture

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

| Тип потоку | Фіксує | Найкраще для |
|---|---|---|
| Standard | Усі типи таблиць і подань та всі DML-зміни — вставки, оновлення, видалення | Таблиці, де може змінитися будь-який рядок (напр. shipments) |
| Append-only | Усі типи таблиць і подань, крім external tables — лише вставки рядків | Таблиці «вставка-один-раз» (напр. delivery events) — ефективніше |
| Insert-only | Зовнішньо керовані Apache Iceberg і external tables — лише вставки рядків | Зовнішні таблиці |
Directory tables показують метадані файлів для сховища (назва, розмір, мітка часу останньої зміни)
Стандартний потік для таблиці 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;
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